如何获取Airflow DynamoDBToS3Operator的导出响应及ExportId
解决DynamoDBToS3Operator获取ExportId的问题
默认DynamoDBToS3Operator执行完成后不会将导出响应推送至XCom,导致下游任务无法获取ExportId及ExportManifest路径。以下是三种可行解决方案:
方案1:自定义Operator推送XCom
重写原Operator的execute方法,将导出响应推送到XCom,下游任务即可正常拉取:
from airflow.providers.amazon.aws.operators.dynamodb_to_s3 import DynamoDBToS3Operator class CustomDynamoDBToS3Operator(DynamoDBToS3Operator): def execute(self, context): # 调用原Operator的执行逻辑 response = super().execute(context) # 将响应推送到XCom context['task_instance'].xcom_push(key='dynamo_export_response', value=response) return response
替换原Operator实例:
dynamo_full_export_task = CustomDynamoDBToS3Operator( task_id="dynamo_full_export_task_task_id", dynamodb_table_name=dag_run_config["dynamo_table"], s3_bucket_name=dag_run_config["bucket"], export_time=pendulum.now("UTC"), s3_key_prefix=dag_run_config["dynamo_export_s3_loc"], export_format="ION", file_size=10**9 )
下游任务拉取并解析:
@task() def extract_dynamo_response(): context = get_current_context() # 从XCom拉取导出响应 response_body = context['task_instance'].xcom_pull( task_ids='dynamo_full_export_task_task_id', key='dynamo_export_response' ) return response_body["ExportManifest"] dynamo_response = extract_dynamo_response()
方案2:直接使用DynamoDBHook调用API
跳过Operator,直接用DynamoDBHook调用导出接口,响应可直接返回给下游任务:
from airflow.providers.amazon.aws.hooks.dynamodb import DynamoDBHook from airflow.decorators import task import pendulum @task() def dynamo_full_export_task(dynamo_table, bucket, s3_prefix): hook = DynamoDBHook(aws_conn_id='aws_default') client = hook.get_conn() # 调用导出API response = client.export_table_to_point_in_time( TableName=dynamo_table, S3Bucket=bucket, S3Prefix=s3_prefix, ExportTime=pendulum.now("UTC"), ExportFormat="ION", FileSizeLimit=10**9 ) # 返回ExportManifest路径 return response["ExportManifest"] # 调用任务并传递参数 dynamo_response = dynamo_full_export_task( dynamo_table=dag_run_config["dynamo_table"], bucket=dag_run_config["bucket"], s3_prefix=dag_run_config["dynamo_export_s3_loc"] )
方案3:通过AWS API查询最近导出记录
如果无法修改Operator或使用Hook,可在下游任务调用list_exports API筛选目标表的最新导出记录:
from airflow.providers.amazon.aws.hooks.dynamodb import DynamoDBHook from airflow.decorators import task import pendulum @task() def get_latest_export_id(dynamo_table, s3_prefix): hook = DynamoDBHook(aws_conn_id='aws_default') client = hook.get_conn() # 查询目标表的导出记录 response = client.list_exports( TableName=dynamo_table, MaxResults=1 ) # 按导出时间排序,取最新的记录 latest_export = sorted( response['ExportSummaries'], key=lambda x: x['ExportTime'], reverse=True )[0] export_id = latest_export['ExportId'] # 拼接ExportManifest完整路径 return f"{s3_prefix}/{export_id}/manifest_summary.json" # 调用任务 dynamo_response = get_latest_export_id( dynamo_table=dag_run_config["dynamo_table"], s3_prefix=dag_run_config["dynamo_export_s3_loc"] )
内容的提问来源于stack exchange,提问作者DataNeer
相关产品推荐
相关产品推荐

