
ingesting-into-data-lake
热门将数据从S3文件、本地上传、JDBC数据库(Oracle、SQL Server、PostgreSQL、MySQL、RDS、Aurora)、Amazon Redshift、Snowflake、BigQuery、DynamoDB或现有Glue目录表(迁移)导入到AWS数据湖。默认目标为S3 Tables;在未采用S3 Tables的环境下,支持在通用桶上使用标准Iceberg。处理一次性加载、周期性管道、迁移。触发词:导入数据、加载数据、摄取、同步数据库、迁移表、将数据移至AWS、设置管道、ETL、从Snowflake拉取、将BigQuery查询结果存入S3、导出DynamoDB、CTAS、转换为Iceberg。请勿用于设置或排查Glue连接(请使用connecting-to-data-source)、创建空表(请使用creating-data-lake-table)、运行查询(请使用querying-data-lake)、通过模糊名称查找表(请使用finding-data-lake-assets)、目录审计(请使用exploring-data-catalog)或SaaS平台(如Salesforce、ServiceNow、SAP、MongoDB、Kafka)。
Import data into the AWS data lake from S3 files, local uploads, JDBC databases (Oracle, SQL Server, PostgreSQL, MySQL, RDS, Aurora), Amazon Redshift, Snowflake, BigQuery, DynamoDB, or existing Glue catalog tables (migration). Default target is S3 Tables; standard Iceberg on a general purpose bucket is supported where S3 Tables is not adopted. Handles one-time loads, recurring pipelines, migrations. Triggers on: import data, load data, ingest, sync database, migrate table, move data to AWS, set up pipeline, ETL, pull from Snowflake, query BigQuery into S3, export DynamoDB, CTAS, convert to Iceberg. Do NOT use for setting up or troubleshooting Glue connections (use connecting-to-data-source), creating empty tables (use creating-data-lake-table), running queries (use querying-data-lake), finding tables by fuzzy name (use finding-data-lake-assets), catalog audit (use exploring-data-catalog), or SaaS platforms like Salesforce, ServiceNow, SAP, MongoDB, Kafka.
摄取到数据湖
将数据从源移动到数据湖中的可查询表。此技能假定源连接(如果需要)已存在。如需设置或排查Glue连接,请委托给connecting-to-data-source。
理念
除非环境另有说明,否则默认使用S3 Tables。 S3 Tables是新的数据湖工作的推荐目标。如果用户的目录清单显示他们尚未采用S3 Tables,则建议在其现有通用桶上使用标准Iceberg,而不是强制他们改变策略。
常见任务
当连接到AWS MCP服务器工具时,您必须使用它们执行命令——它们提供验证、沙盒执行和审计日志记录。仅在MCP不可用时回退到AWS CLI。您必须在执行前解释每个步骤。
工作流程
1. 验证依赖项和上下文
- 您必须检查AWS MCP工具或AWS CLI是否可用,如果缺失则告知用户
- 您必须确认目标AWS区域,并使用
aws sts get-caller-identity验证凭证 - 对于SageMaker Unified Studio项目角色,请注意目标表和连接可能限定在项目范围内。请参阅
querying-data-lake中的调用者ARN检测模式。
2. 分类源
| 用户说... | 源类型 | 参考 |
|---|---|---|
| "上传我的文件", "本地CSV", "移动到S3" | 本地文件 | local-upload.md |
| "从S3加载", "从s3://导入CSV/JSON/Parquet" | S3文件 | s3-files.md |
| "从Oracle/Postgres/MySQL/SQL Server/Redshift/RDS/Aurora导入" | JDBC | jdbc-ingest.md |
| "从Snowflake拉取", "Snowflake表到S3" | Snowflake | snowflake-ingest.md |
| "从BigQuery导入", "GCP分析到S3" | BigQuery | bigquery-ingest.md |
| "导出DynamoDB", "DynamoDB到数据湖" | DynamoDB | dynamodb-ingest.md |
| "迁移Glue表", "将Hive转换为Iceberg" | 目录迁移 | catalog-migration.md |
如果用户提到Salesforce、ServiceNow、SAP、MongoDB、Kafka或其他SaaS/流式源,请拒绝——这些在此版本中不受支持。
如果源表通过模糊名称或业务名称引用("迁移我们的订单表", "从销售仓库拉取"),请先委托给finding-data-lake-assets进行解析,然后再继续。
3. 确认连接存在(如适用)
对于JDBC、Snowflake和BigQuery源,需要Glue连接。检查:
aws glue get-connection --name <CONNECTION_NAME> --region <REGION>
如果连接不存在,请停止并委托给connecting-to-data-source创建和测试。在连接验证通过之前,不要继续摄取。
本地文件、S3文件、DynamoDB和目录迁移不需要Glue连接。
4. 明确目标
在创建或写入任何表之前,您必须询问用户(或根据目录清单建议):
- 数据库/命名空间:是否存在特定的目标数据库?还是应该创建一个?
- 表:现有表(追加/合并)还是新表(委托给
creating-data-lake-table)? - 格式:S3 Tables(默认)、标准Iceberg还是原始Parquet?
清单感知的默认值:
如果您已经运行过exploring-data-catalog或可以快速检查,请使用已有的:
- 账户有
s3tablescatalog联合目录和活跃的表桶:推荐S3 Tables - 账户有通用桶且包含Iceberg表,但未使用S3 Tables:推荐在其现有桶上使用标准Iceberg
- 账户在S3上使用Parquet/ORC且没有Iceberg元数据:询问是否现在采用Iceberg(推荐是)或继续使用原始文件
不要强制未采用S3 Tables的客户使用S3 Tables。请参阅iceberg-catalog-config-and-usage.md。
此步骤的委托:
- 目标表不存在 ->
creating-data-lake-table - 目标数据库通过模糊术语命名 ->
finding-data-lake-assets - 用户不知道存在什么 ->
exploring-data-catalog
5. 执行源工作流
阅读源特定的参考并遵循其阶段。每个参考都包含作业模板、注意事项和故障排除:
- 本地 / S3 / JDBC / Snowflake / BigQuery / DynamoDB / 目录迁移——每个源一个参考
常见的Glue 5.1或更高版本作业配置和PySpark模板共享在glue-job-config.md和glue-job-scripts.md中。
6. 验证
运行以下所有三项,不要跳过:
- 行数与预期匹配(源 vs 目标)
- 关键列的空值检查
- 抽查3-5行样本数据
请参阅data-quality-validation.md。
7. 调度(如果周期性)
对于周期性管道,使用cron调度创建Glue触发器。请参阅testing-and-scheduling.md。简单的单步管道使用Glue触发器;多步带分支的管道使用MWAA。
参数路由
- 仅S3路径:推断为一次性加载,从步骤2开始处理S3文件
- 连接名称:从步骤3开始,使用指定的连接
- 表名称:从步骤4开始,询问这是源还是目标
--target标志:在步骤4中预填目标格式- 无参数:交互式引导
注意事项
- S3 Tables需要Glue 5.1或更高版本以及
--datalake-formats iceberg作业参数 - 所有
spark.sql.catalog.*配置必须放在--conf作业参数中,绝不能放在spark.conf.set()中。否则Glue 5.x会抛出AnalysisException: Cannot modify the value of a static config。请参阅iceberg-catalog-config-and-usage.md了解正确的目录配置。 - S3 Tables目录配置中需要
warehouse参数。没有它,Spark会失败并显示"Cannot derive default warehouse location"。 - S3 Tables中的表和列名称必须全部小写
overwritePartitions()仅替换DataFrame中存在的分区——对于需要删除的完全刷新,请使用createOrReplace()- 标准Iceberg目标必须包含LOCATION子句;S3 Tables则不能包含
- DynamoDB不需要Glue连接——不要尝试创建
- 摄取期间的连接失败委托给
connecting-to-data-source;不要在此技能中调试网络/凭证 - 对于SageMaker Unified Studio项目中的目标表,确保在Glue作业运行之前项目角色对目标命名空间具有写权限
故障排除
| 错误 | 可能原因 | 操作 |
|---|---|---|
| S3访问被拒绝 | 缺少IAM权限 | 检查Glue角色是否具有s3:GetObject、s3:PutObject |
| S3 Tables访问被拒绝 | 缺少s3tables:*权限 | 向Glue角色添加S3 Tables内联策略 |
| CTAS超时 | 数据集对于Athena过大 | 切换到Glue ETL或使用WHERE过滤器分批处理 |
| JDBC连接超时/认证失败 | 连接级别问题 | 委托给connecting-to-data-source |
| DynamoDB吞吐量超出 | 读取百分比过高 | 降低read.percent或使用原生导出 |
请参阅error-handling.md获取完整目录。
参考
源特定
- local-upload.md -- 本地文件
- s3-files.md -- S3文件(CSV、JSON、Parquet、Avro、ORC)
- jdbc-ingest.md -- Oracle、SQL Server、PostgreSQL、MySQL、RDS、Aurora、Redshift
- snowflake-ingest.md -- Snowflake
- bigquery-ingest.md -- BigQuery
- dynamodb-ingest.md -- DynamoDB(导出和Glue直接读取)
- catalog-migration.md -- 现有Glue目录表(Hive、自管理Iceberg)
交叉领域
- iceberg-catalog-config-and-usage.md -- S3 Tables、标准Iceberg、原始文件:目录配置、引擎访问模式
- glue-job-config.md -- 作业大小、监控、重试
- glue-job-scripts.md -- PySpark模板(追加、更新插入、自定义SQL、完全刷新)
- incremental-loading.md -- 水印策略
- testing-and-scheduling.md -- Glue触发器、MWAA
- data-quality-validation.md -- 行数、空值检查、Glue数据质量
- schema-evolution.md -- ALTER TABLE ADD COLUMNS、嵌套JSON
- type-transformations.md -- 类型冲突解决
- format-specific-loading.md -- CSV/JSON/Parquet/Avro/ORC细节
- athena-loading.md -- Athena INSERT INTO作为简单加载回退
- error-handling.md -- 摄取错误(连接错误委托给connecting-to-data-source)
- upload-options.md -- aws s3 cp vs sync、多部分
迁移特定
- ctas-patterns.md -- Athena CTAS语法和分区转换
- glue-etl-migration.md -- 通过Glue 5.1或更高版本PySpark进行大表迁移
- migration-validation.md -- 完整验证清单
- migration-troubleshooting.md -- CTAS失败、可见性、分区
JDBC特定
- jdbc-schema-discovery.md -- 爬网程序、直接检查、自定义SQL
- jdbc-performance.md -- 并行读取、分区





