Current section

Files

Jump to
aws_erlang src aws_data_pipeline.erl
Raw

src/aws_data_pipeline.erl

%% WARNING: DO NOT EDIT, AUTO-GENERATED CODE!
%% See https://github.com/aws-beam/aws-codegen for more details.
%% @doc AWS Data Pipeline configures and manages a data-driven workflow
%% called a pipeline.
%%
%% AWS Data Pipeline handles the details of scheduling and ensuring that data
%% dependencies are met so that your application can focus on processing the
%% data.
%%
%% AWS Data Pipeline provides a JAR implementation of a task runner called
%% AWS Data Pipeline Task Runner. AWS Data Pipeline Task Runner provides
%% logic for common data management scenarios, such as performing database
%% queries and running data analysis using Amazon Elastic MapReduce (Amazon
%% EMR). You can use AWS Data Pipeline Task Runner as your task runner, or
%% you can write your own task runner to provide custom data management.
%%
%% AWS Data Pipeline implements two main sets of functionality. Use the first
%% set to create a pipeline and define data sources, schedules, dependencies,
%% and the transforms to be performed on the data. Use the second set in your
%% task runner application to receive the next task ready for processing. The
%% logic for performing the task, such as querying the data, running data
%% analysis, or converting the data from one format to another, is contained
%% within the task runner. The task runner performs the task assigned to it
%% by the web service, reporting progress to the web service as it does so.
%% When the task is done, the task runner reports the final success or
%% failure of the task to the web service.
-module(aws_data_pipeline).
-export([activate_pipeline/2,
activate_pipeline/3,
add_tags/2,
add_tags/3,
create_pipeline/2,
create_pipeline/3,
deactivate_pipeline/2,
deactivate_pipeline/3,
delete_pipeline/2,
delete_pipeline/3,
describe_objects/2,
describe_objects/3,
describe_pipelines/2,
describe_pipelines/3,
evaluate_expression/2,
evaluate_expression/3,
get_pipeline_definition/2,
get_pipeline_definition/3,
list_pipelines/2,
list_pipelines/3,
poll_for_task/2,
poll_for_task/3,
put_pipeline_definition/2,
put_pipeline_definition/3,
query_objects/2,
query_objects/3,
remove_tags/2,
remove_tags/3,
report_task_progress/2,
report_task_progress/3,
report_task_runner_heartbeat/2,
report_task_runner_heartbeat/3,
set_status/2,
set_status/3,
set_task_status/2,
set_task_status/3,
validate_pipeline_definition/2,
validate_pipeline_definition/3]).
-include_lib("hackney/include/hackney_lib.hrl").
%%====================================================================
%% API
%%====================================================================
%% @doc Validates the specified pipeline and starts processing pipeline
%% tasks.
%%
%% If the pipeline does not pass validation, activation fails.
%%
%% If you need to pause the pipeline to investigate an issue with a
%% component, such as a data source or script, call `DeactivatePipeline'.
%%
%% To activate a finished pipeline, modify the end date for the pipeline and
%% then activate it.
activate_pipeline(Client, Input)
when is_map(Client), is_map(Input) ->
activate_pipeline(Client, Input, []).
activate_pipeline(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"ActivatePipeline">>, Input, Options).
%% @doc Adds or modifies tags for the specified pipeline.
add_tags(Client, Input)
when is_map(Client), is_map(Input) ->
add_tags(Client, Input, []).
add_tags(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"AddTags">>, Input, Options).
%% @doc Creates a new, empty pipeline.
%%
%% Use `PutPipelineDefinition' to populate the pipeline.
create_pipeline(Client, Input)
when is_map(Client), is_map(Input) ->
create_pipeline(Client, Input, []).
create_pipeline(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"CreatePipeline">>, Input, Options).
%% @doc Deactivates the specified running pipeline.
%%
%% The pipeline is set to the `DEACTIVATING' state until the deactivation
%% process completes.
%%
%% To resume a deactivated pipeline, use `ActivatePipeline'. By default,
%% the pipeline resumes from the last completed execution. Optionally, you
%% can specify the date and time to resume the pipeline.
deactivate_pipeline(Client, Input)
when is_map(Client), is_map(Input) ->
deactivate_pipeline(Client, Input, []).
deactivate_pipeline(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"DeactivatePipeline">>, Input, Options).
%% @doc Deletes a pipeline, its pipeline definition, and its run history.
%%
%% AWS Data Pipeline attempts to cancel instances associated with the
%% pipeline that are currently being processed by task runners.
%%
%% Deleting a pipeline cannot be undone. You cannot query or restore a
%% deleted pipeline. To temporarily pause a pipeline instead of deleting it,
%% call `SetStatus' with the status set to `PAUSE' on individual
%% components. Components that are paused by `SetStatus' can be resumed.
delete_pipeline(Client, Input)
when is_map(Client), is_map(Input) ->
delete_pipeline(Client, Input, []).
delete_pipeline(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"DeletePipeline">>, Input, Options).
%% @doc Gets the object definitions for a set of objects associated with the
%% pipeline.
%%
%% Object definitions are composed of a set of fields that define the
%% properties of the object.
describe_objects(Client, Input)
when is_map(Client), is_map(Input) ->
describe_objects(Client, Input, []).
describe_objects(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"DescribeObjects">>, Input, Options).
%% @doc Retrieves metadata about one or more pipelines.
%%
%% The information retrieved includes the name of the pipeline, the pipeline
%% identifier, its current state, and the user account that owns the
%% pipeline. Using account credentials, you can retrieve metadata about
%% pipelines that you or your IAM users have created. If you are using an IAM
%% user account, you can retrieve metadata about only those pipelines for
%% which you have read permissions.
%%
%% To retrieve the full pipeline definition instead of metadata about the
%% pipeline, call `GetPipelineDefinition'.
describe_pipelines(Client, Input)
when is_map(Client), is_map(Input) ->
describe_pipelines(Client, Input, []).
describe_pipelines(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"DescribePipelines">>, Input, Options).
%% @doc Task runners call `EvaluateExpression' to evaluate a string in
%% the context of the specified object.
%%
%% For example, a task runner can evaluate SQL queries stored in Amazon S3.
evaluate_expression(Client, Input)
when is_map(Client), is_map(Input) ->
evaluate_expression(Client, Input, []).
evaluate_expression(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"EvaluateExpression">>, Input, Options).
%% @doc Gets the definition of the specified pipeline.
%%
%% You can call `GetPipelineDefinition' to retrieve the pipeline
%% definition that you provided using `PutPipelineDefinition'.
get_pipeline_definition(Client, Input)
when is_map(Client), is_map(Input) ->
get_pipeline_definition(Client, Input, []).
get_pipeline_definition(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"GetPipelineDefinition">>, Input, Options).
%% @doc Lists the pipeline identifiers for all active pipelines that you have
%% permission to access.
list_pipelines(Client, Input)
when is_map(Client), is_map(Input) ->
list_pipelines(Client, Input, []).
list_pipelines(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"ListPipelines">>, Input, Options).
%% @doc Task runners call `PollForTask' to receive a task to perform from
%% AWS Data Pipeline.
%%
%% The task runner specifies which tasks it can perform by setting a value
%% for the `workerGroup' parameter. The task returned can come from any
%% of the pipelines that match the `workerGroup' value passed in by the
%% task runner and that was launched using the IAM user credentials specified
%% by the task runner.
%%
%% If tasks are ready in the work queue, `PollForTask' returns a response
%% immediately. If no tasks are available in the queue, `PollForTask'
%% uses long-polling and holds on to a poll connection for up to a 90
%% seconds, during which time the first newly scheduled task is handed to the
%% task runner. To accomodate this, set the socket timeout in your task
%% runner to 90 seconds. The task runner should not call `PollForTask'
%% again on the same `workerGroup' until it receives a response, and this
%% can take up to 90 seconds.
poll_for_task(Client, Input)
when is_map(Client), is_map(Input) ->
poll_for_task(Client, Input, []).
poll_for_task(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"PollForTask">>, Input, Options).
%% @doc Adds tasks, schedules, and preconditions to the specified pipeline.
%%
%% You can use `PutPipelineDefinition' to populate a new pipeline.
%%
%% `PutPipelineDefinition' also validates the configuration as it adds it
%% to the pipeline. Changes to the pipeline are saved unless one of the
%% following three validation errors exists in the pipeline.
%%
%% <ol> <li>An object is missing a name or identifier field.</li> <li>A
%% string or reference field is empty.</li> <li>The number of objects in the
%% pipeline exceeds the maximum allowed objects.</li> <li>The pipeline is in
%% a FINISHED state.</li> </ol> Pipeline object definitions are passed to the
%% `PutPipelineDefinition' action and returned by the
%% `GetPipelineDefinition' action.
put_pipeline_definition(Client, Input)
when is_map(Client), is_map(Input) ->
put_pipeline_definition(Client, Input, []).
put_pipeline_definition(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"PutPipelineDefinition">>, Input, Options).
%% @doc Queries the specified pipeline for the names of objects that match
%% the specified set of conditions.
query_objects(Client, Input)
when is_map(Client), is_map(Input) ->
query_objects(Client, Input, []).
query_objects(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"QueryObjects">>, Input, Options).
%% @doc Removes existing tags from the specified pipeline.
remove_tags(Client, Input)
when is_map(Client), is_map(Input) ->
remove_tags(Client, Input, []).
remove_tags(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"RemoveTags">>, Input, Options).
%% @doc Task runners call `ReportTaskProgress' when assigned a task to
%% acknowledge that it has the task.
%%
%% If the web service does not receive this acknowledgement within 2 minutes,
%% it assigns the task in a subsequent `PollForTask' call. After this
%% initial acknowledgement, the task runner only needs to report progress
%% every 15 minutes to maintain its ownership of the task. You can change
%% this reporting time from 15 minutes by specifying a
%% `reportProgressTimeout' field in your pipeline.
%%
%% If a task runner does not report its status after 5 minutes, AWS Data
%% Pipeline assumes that the task runner is unable to process the task and
%% reassigns the task in a subsequent response to `PollForTask'. Task
%% runners should call `ReportTaskProgress' every 60 seconds.
report_task_progress(Client, Input)
when is_map(Client), is_map(Input) ->
report_task_progress(Client, Input, []).
report_task_progress(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"ReportTaskProgress">>, Input, Options).
%% @doc Task runners call `ReportTaskRunnerHeartbeat' every 15 minutes to
%% indicate that they are operational.
%%
%% If the AWS Data Pipeline Task Runner is launched on a resource managed by
%% AWS Data Pipeline, the web service can use this call to detect when the
%% task runner application has failed and restart a new instance.
report_task_runner_heartbeat(Client, Input)
when is_map(Client), is_map(Input) ->
report_task_runner_heartbeat(Client, Input, []).
report_task_runner_heartbeat(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"ReportTaskRunnerHeartbeat">>, Input, Options).
%% @doc Requests that the status of the specified physical or logical
%% pipeline objects be updated in the specified pipeline.
%%
%% This update might not occur immediately, but is eventually consistent. The
%% status that can be set depends on the type of object (for example,
%% DataNode or Activity). You cannot perform this operation on `FINISHED'
%% pipelines and attempting to do so returns `InvalidRequestException'.
set_status(Client, Input)
when is_map(Client), is_map(Input) ->
set_status(Client, Input, []).
set_status(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"SetStatus">>, Input, Options).
%% @doc Task runners call `SetTaskStatus' to notify AWS Data Pipeline
%% that a task is completed and provide information about the final status.
%%
%% A task runner makes this call regardless of whether the task was
%% sucessful. A task runner does not need to call `SetTaskStatus' for
%% tasks that are canceled by the web service during a call to
%% `ReportTaskProgress'.
set_task_status(Client, Input)
when is_map(Client), is_map(Input) ->
set_task_status(Client, Input, []).
set_task_status(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"SetTaskStatus">>, Input, Options).
%% @doc Validates the specified pipeline definition to ensure that it is well
%% formed and can be run without error.
validate_pipeline_definition(Client, Input)
when is_map(Client), is_map(Input) ->
validate_pipeline_definition(Client, Input, []).
validate_pipeline_definition(Client, Input, Options)
when is_map(Client), is_map(Input), is_list(Options) ->
request(Client, <<"ValidatePipelineDefinition">>, Input, Options).
%%====================================================================
%% Internal functions
%%====================================================================
-spec request(aws_client:aws_client(), binary(), map(), list()) ->
{ok, Result, {integer(), list(), hackney:client()}} |
{error, Error, {integer(), list(), hackney:client()}} |
{error, term()} when
Result :: map() | undefined,
Error :: map().
request(Client, Action, Input, Options) ->
RequestFun = fun() -> do_request(Client, Action, Input, Options) end,
aws_request:request(RequestFun, Options).
do_request(Client, Action, Input0, Options) ->
Client1 = Client#{service => <<"datapipeline">>},
Host = build_host(<<"datapipeline">>, Client1),
URL = build_url(Host, Client1),
Headers = [
{<<"Host">>, Host},
{<<"Content-Type">>, <<"application/x-amz-json-1.1">>},
{<<"X-Amz-Target">>, <<"DataPipeline.", Action/binary>>}
],
Input = Input0,
Payload = jsx:encode(Input),
SignedHeaders = aws_request:sign_request(Client1, <<"POST">>, URL, Headers, Payload),
Response = hackney:request(post, URL, SignedHeaders, Payload, Options),
handle_response(Response).
handle_response({ok, 200, ResponseHeaders, Client}) ->
case hackney:body(Client) of
{ok, <<>>} ->
{ok, undefined, {200, ResponseHeaders, Client}};
{ok, Body} ->
Result = jsx:decode(Body),
{ok, Result, {200, ResponseHeaders, Client}}
end;
handle_response({ok, StatusCode, ResponseHeaders, Client}) ->
{ok, Body} = hackney:body(Client),
Error = jsx:decode(Body),
{error, Error, {StatusCode, ResponseHeaders, Client}};
handle_response({error, Reason}) ->
{error, Reason}.
build_host(_EndpointPrefix, #{region := <<"local">>, endpoint := Endpoint}) ->
Endpoint;
build_host(_EndpointPrefix, #{region := <<"local">>}) ->
<<"localhost">>;
build_host(EndpointPrefix, #{region := Region, endpoint := Endpoint}) ->
aws_util:binary_join([EndpointPrefix, Region, Endpoint], <<".">>).
build_url(Host, Client) ->
Proto = aws_client:proto(Client),
Port = aws_client:port(Client),
aws_util:binary_join([Proto, <<"://">>, Host, <<":">>, Port, <<"/">>], <<"">>).