免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

Apache Airflow 实战:使用 S3ToDynamoDBOperator 将 S3 数据导入 DynamoDB(新表与现有表两种模式)

Apache Airflow 实战:使用 S3ToDynamoDBOperator 将 S3 数据导入 DynamoDB(新表与现有表两种模式) Apache Airflow 实战使用 S3ToDynamoDBOperator 将 S3 数据导入 DynamoDB新表与现有表两种模式【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本指南基于 Apache Airflow Amazon 提供商包apache-airflow-providers-amazon中的传输算子S3ToDynamoDBOperator系统讲解如何把存放在 Amazon S3 桶中的数据CSV、DynamoDB JSON 或 Amazon ION 格式加载到新建的或已存在的 Amazon DynamoDB 表。读完本文你将掌握算子全部核心参数的含义与默认值、导入新表与导入现有表两条执行路径的底层机制、错误处理与容错选项并能够参照仓库自带的系统测试示例example_s3_to_dynamodb.py快速落地一个可运行的传输 DAG。一、功能概述一条命令打通 S3 → DynamoDBS3ToDynamoDBOperator是 Amazon 提供商包提供的 transfer 类算子它的职责非常聚焦将 S3 桶中存储的数据加载到一个新建的或已存在的 DynamoDB 表。它底层使用 Amazon DynamoDB 的ImportTable服务该服务在导入过程中会与 Amazon S3、CloudWatch 等多个 AWS 服务协同工作。其完整实现位于 s3_to_dynamodb.py算子的类定义S3ToDynamoDBOperator(BaseOperator)以类型化字典约束了导入所需的表结构定义AttributeDefinition定义属性名与类型AttributeType取值限定为S字符串、N数字、B二进制KeySchema定义主键与排序键KeyType取值限定为HASH分区键或RANGE排序键。二、前置条件在正式使用本算子之前需要完成三件事详见 prerequisite_tasks.rst准备 AWS 资源使用 AWS Console 或 AWS CLI 预先创建好所需的 S3 桶、目标 DynamoDB 表若走“导入现有表”路径等资源并确保 IAM 权限允许 DynamoDB 的ImportTable、Scan、BatchWriteItem、DeleteTable等操作以及对应的 S3 读取权限。安装依赖库通过 pip 安装 Amazon 提供商包pip install apache-airflow[amazon]更详细的安装说明可参考 airflow-core 安装文档。配置 AWS 连接在 Airflow 中配置好 AWS 连接Connection算子默认使用连接 IDaws_default配置方法见 AWS 连接配置。三、导入到新建 DynamoDB 表默认行为3.1 示例 DAG仓库的系统测试示例 DAG example_s3_to_dynamodb.py 提供了一个完整可参考的最小示例。测试构造了一个以cocktail_id为 HASH 主键的表结构并准备了如下 CSV 样例数据cocktail_id,cocktail_name,base_spirit 1,Caipirinha,Cachaca 2,Bramble,Gin 3,Daiquiri,Rum将 S3 数据导入新表的核心任务定义如下对应原文档[START howto_transfer_s3_to_dynamodb]片段transfer_1 S3ToDynamoDBOperator( task_ids3_to_dynamodb, s3_bucketbucket_name, s3_keys3_key, dynamodb_table_namenew_table_name, input_formatCSV, import_table_kwargs{ InputFormatOptions: { Csv: { Delimiter: ,, } } }, dynamodb_attributes[ {AttributeName: cocktail_id, AttributeType: S}, ], dynamodb_key_schema[ {AttributeName: cocktail_id, KeyType: HASH}, ], )3.2 底层执行流程从源码execute()与_load_into_new_table()的实现s3_to_dynamodb.py可以看到导入新表路径做了以下事情通过DynamoDBHook见 dynamodb.py获取 boto3 客户端调用client.import_table(...)请求体结构如下response client.import_table( S3BucketSource{ S3Bucket: self.s3_bucket, S3KeyPrefix: self.s3_key, }, InputFormatself.input_format, TableCreationParameters{ TableName: table_name, AttributeDefinitions: self.dynamodb_attributes, KeySchema: self.dynamodb_key_schema, BillingMode: self.billing_mode, **import_table_creation_config, }, **import_table_config, )若ImportTableDescription.ImportStatus直接为FAILED立即抛出AirflowException包含FailureCode与FailureMessage例如测试中的invalid csv format见 test_s3_to_dynamodb.py若wait_for_completionTrue默认则通过 DynamoDB 的import_tablewaiter 轮询导入任务直到完成执行成功时返回导入任务的ImportArn即任务返回值可被下游任务通过 XCom 引用。从单元测试test_s3_to_dynamodb_new_table_wait_for_completion的断言可以看到该调用细节import_table被调用一次、waiter 等待时传入WaiterConfig{Delay: 30, MaxAttempts: 240}返回值正是ImportArn。四、导入到现有 DynamoDB 表自定义实现路径4.1 为什么需要特殊处理DynamoDB 的 ImportTable 服务目前不支持直接导入到已存在的表。因此算子采用了一种自定义的两阶段方案源码第 207-243 行创建一个临时 DynamoDB 表名称规则为{dynamodb_tmp_table_prefix}_{dynamodb_table_name}前缀默认tmp调用_load_into_new_table让 ImportTable 服务先把 S3 数据批量加载进临时表使用 boto3 的scan分页器以ConsistentReadTrue、SelectALL_ATTRIBUTES全量扫描临时表将记录按页取回内存通过DynamoDBHook.write_batch_data使用Table.batch_writer将每条记录put_item写入目标表overwrite_by_pkeys传入主键列表实现按主键覆盖写入无论成功与否在finally块中删除临时表。4.2 示例代码对应用原文档[START howto_transfer_s3_to_dynamodb_existing_table]片段只需把use_existing_table设为True即可transfer_2 S3ToDynamoDBOperator( task_ids3_to_dynamodb_new_table, s3_bucketbucket_name, s3_keys3_key, dynamodb_table_nameexisting_table_name, use_existing_tableTrue, input_formatCSV, import_table_kwargs{ InputFormatOptions: { Csv: { Delimiter: ,, } } }, dynamodb_attributes[ {AttributeName: cocktail_id, AttributeType: S}, ], dynamodb_key_schema[ {AttributeName: cocktail_id, KeyType: HASH}, ], )4.3 重要限制源码_load_into_existing_table()开头有一处硬性校验if not self.wait_for_completion: raise ValueError(wait_for_completion must be set to True when loading into an existing table)即导入现有表时wait_for_completion必须保持为True否则会抛出ValueError。原因很直观后续的 scan 操作依赖临时表导入完成。单元测试test_s3_to_dynamodb_existing_tabletest_s3_to_dynamodb.py完整验证了这条路径先以tmp_test-table为名调用_load_into_new_table然后get_paginator(scan)分页扫描、batch_writer(overwrite_by_pkeys[attribute_a])批量写入目标表test-table最后delete_table删除临时表返回值为目标表 ARN。五、核心参数详解以下参数均来自算子__init__签名s3_to_dynamodb.py参数必填默认值说明s3_bucket是—被导入的 S3 桶名称s3_key是—S3 键或键前缀可匹配单个或多个对象支持前缀批量导入dynamodb_table_name是—目标表名称导入现有表时即为现有表名dynamodb_key_schema是—主键与排序键列表KeyType取值HASH/RANGEdynamodb_attributes否None表的属性定义列表AttributeType取值S/N/Bdynamodb_tmp_table_prefix否tmp临时表名前缀仅导入现有表时使用delete_on_error否False导入出错时是否删除新建的或临时的DynamoDB 表use_existing_table否False是否导入到已存在的表True走“临时表 scan 批量写”路径input_format否DYNAMODB_JSON导入数据格式取值CSV/DYNAMODB_JSON/IONbilling_mode否PAY_PER_REQUEST表计费模式取值PROVISIONED/PAY_PER_REQUESTimport_table_kwargs否None透传给import_table的附加参数如ClientToken、InputCompressionType、InputFormatOptionsimport_table_creation_kwargs否None透传给TableCreationParameters的附加参数如ProvisionedThroughput、SSESpecification、GlobalSecondaryIndexeswait_for_completion否True是否等待导入任务完成导入现有表时必须为Truecheck_interval否30每次状态检查的间隔秒数waiter 的Delaymax_attempts否240完成检查的最大尝试次数waiter 的MaxAttemptsaws_conn_id否aws_default使用的 AWS 连接 ID两点使用建议指定 CSV 分隔符当input_formatCSV时示例通过import_table_kwargs{InputFormatOptions: {Csv: {Delimiter: ,}}}显式声明分隔符如果你的数据是制表符或分号分隔在这里对应调整。切换计费模式默认PAY_PER_REQUEST按请求计费若希望使用预留吞吐可设置billing_modePROVISIONED并通过import_table_creation_kwargs{ProvisionedThroughput: {...}}传入读写容量。六、错误处理与容错机制从源码与单元测试可以归纳出算子的完整错误处理链调用阶段失败import_table抛出ClientError时记录错误日志并抛出AirflowException(S3 load into DynamoDB table failed with error: ...)对应测试test_s3_to_dynamodb_new_table_client_error任务创建即失败响应中ImportStatus FAILED抛出包含FailureCode/FailureMessage的AirflowException对应测试test_s3_to_dynamodb_new_table_job_startup_error等待期间失败WaiterError先通过DynamoDBHook.get_import_status内部调用describe_import获取最新状态、错误码与错误信息再抛出AirflowExceptiondelete_on_error联动在上述第 3 种场景中若delete_on_errorTrue算子会先调用client.delete_table(TableNametable_name)清理半成品表再抛异常对应测试test_s3_to_dynamodb_new_table_delete_on_error的两种参数化场景delete-on-error与no-delete-on-error。此外DynamoDBHook.get_import_status对ImportNotFoundException会抛出“导入任务未找到”的AirflowException用于兜底异常状态。七、完整 DAG 串接参考系统测试 DAG 展示了完整的“准备数据 → 传输 → 清理”链路可作为生产 DAG 的骨架模板with DAG( dag_idexample_s3_to_dynamodb, scheduleonce, start_datedatetime(2021, 1, 1), catchupFalse, ) as dag: # 准备创建目标表、创建 S3 桶、写入 CSV 对象 create_table set_up_table(table_nameexisting_table_name) create_bucket S3CreateBucketOperator(task_idcreate_bucket, bucket_namebucket_name) create_object S3CreateObjectOperator( task_idcreate_object, s3_bucketbucket_name, s3_keys3_key, dataSAMPLE_DATA, replaceTrue, ) # 传输新建表 现有表两条路径 transfer_1 S3ToDynamoDBOperator(...) # 导入新表 transfer_2 S3ToDynamoDBOperator(...) # use_existing_tableTrue # 清理删除表与桶 delete_existing_table delete_dynamodb_table(table_nameexisting_table_name) delete_new_table delete_dynamodb_table(table_namenew_table_name) delete_bucket S3DeleteBucketOperator( task_iddelete_bucket, bucket_namebucket_name, trigger_ruleTriggerRule.ALL_DONE, force_deleteTrue, ) chain( test_context, create_table, create_bucket, wait_for_bucket(s3_bucket_namebucket_name), create_object, transfer_1, transfer_2, delete_existing_table, delete_new_table, delete_bucket, )其中set_up_table使用 boto3 创建带ProvisionedThroughput{ReadCapacityUnits: 1, WriteCapacityUnits: 1}的目标表delete_dynamodb_table在TriggerRule.ALL_DONE下确保无论任务成败都会清理资源。八、小结S3ToDynamoDBOperator屏蔽了 DynamoDB ImportTable 服务的细节让 Airflow 用户可以用声明式参数完成 S3 → DynamoDB 的数据搬运默认路径走云上 ImportTable 批量导入新表适合大数据量冷加载use_existing_tableTrue路径则通过“临时表 一致性扫描 batch_writer 批量写入”在应用层实现了对现有表的增量灌入代价是数据需要经过内存中转适合中小规模数据。建议结合本文参数表参考 example_s3_to_dynamodb.py 与 test_s3_to_dynamodb.py 中的断言细节验证你的调用参数是否符合预期再投入到生产 DAG 中。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表