如何使用AWS CDK通过EMR Serverless应用创建ETL任务?
用AWS CDK通过EMR Serverless应用创建ETL任务
当然可以通过EMR Serverless应用创建ETL任务,你已经完成了应用创建的步骤,接下来只需要在CDK中定义Job Run(即具体的ETL任务)并关联到已创建的应用即可。以下是具体实现方案:
核心逻辑
EMR Serverless的ETL任务以Job Run的形式存在,需要关联到已启动的Serverless应用。CDK中可以通过CfnJobRun(L1构造)或更高层的L2构造来定义任务,同时结合IAM权限、S3存储和可选的触发规则完成完整的ETL流程。
1. 依赖准备
确保CDK项目中引入了EMR Serverless的依赖包,以TypeScript为例:
npm install @aws-cdk/aws-emrserverless
2. 基础ETL任务定义(一次性任务)
假设你已经通过CDK创建了Spark类型的EMR Serverless应用,以下代码展示如何提交一个Spark ETL任务:
import * as cdk from 'aws-cdk-lib'; import * as emrserverless from '@aws-cdk/aws-emrserverless'; import * as s3 from '@aws-cdk/aws-s3'; export class EmrServerlessEtlStack extends cdk.Stack { constructor(scope: cdk.App, id: string, props?: cdk.StackProps) { super(scope, id, props); // 引用已创建的EMR Serverless应用 const existingSparkApp = emrserverless.CfnApplication.fromCfnApplicationAttributes(this, 'ExistingSparkApp', { applicationId: 'your-app-id', // 替换为你的应用ID arn: 'arn:aws:emr-serverless:us-east-1:123456789012:application/your-app-id' }); // 引用存储ETL脚本、输入输出数据的S3桶 const etlBucket = s3.Bucket.fromBucketName(this, 'EtlBucket', 'your-etl-bucket'); // 创建ETL任务(Job Run) new emrserverless.CfnJobRun(this, 'SparkEtlJob', { applicationId: existingSparkApp.applicationId, executionRoleArn: 'arn:aws:iam::123456789012:role/emr-serverless-exec-role', // 替换为有权限的执行角色ARN jobDriver: { sparkSubmit: { entryPoint: `s3://${etlBucket.bucketName}/scripts/etl-transform.py`, // ETL脚本的S3路径 entryPointArguments: [ `s3://${etlBucket.bucketName}/raw-data/`, `s3://${etlBucket.bucketName}/processed-data/` ], // 传递给脚本的输入输出路径参数 sparkSubmitParameters: '--conf spark.executor.instances=2 --conf spark.executor.memory=4g --conf spark.driver.memory=2g' // Spark资源配置 } }, configurationOverrides: { monitoringConfiguration: { s3MonitoringConfiguration: { logUri: `s3://${etlBucket.bucketName}/job-logs/` // 任务日志存储路径 } } } }); } }
3. 关键配置说明
- executionRoleArn:该IAM角色必须具备以下权限:访问ETL脚本、输入输出S3路径的权限;EMR Serverless任务执行的相关权限;写入CloudWatch日志或S3日志的权限。遵循最小权限原则,避免使用过度宽泛的权限策略。
- jobDriver:根据你的ETL引擎选择配置(Spark/Hive),示例中使用SparkSubmit,
entryPoint指向你的ETL代码,entryPointArguments传递业务参数,sparkSubmitParameters可调整资源分配、Spark运行参数。 - monitoringConfiguration:配置日志存储位置,方便排查任务失败原因、监控运行状态。
4. 定时触发ETL任务
如果需要定期执行ETL(如每日同步),可以结合CloudWatch Events规则实现定时触发:
import * as events from '@aws-cdk/aws-events'; import * as targets from '@aws-cdk/aws-events-targets'; // 创建每日凌晨1点触发的规则 const dailyEtlRule = new events.Rule(this, 'DailyEtlTrigger', { schedule: events.Schedule.cron({ hour: '1', minute: '0' }) }); // 将EMR Serverless任务作为触发目标 dailyEtlRule.addTarget(new targets.EmrServerlessStartJobRun({ applicationId: existingSparkApp.applicationId, executionRoleArn: 'arn:aws:iam::123456789012:role/emr-serverless-exec-role', jobDriver: { sparkSubmit: { entryPoint: `s3://${etlBucket.bucketName}/scripts/etl-transform.py`, entryPointArguments: [ `s3://${etlBucket.bucketName}/raw-data/`, `s3://${etlBucket.bucketName}/processed-data/` ], sparkSubmitParameters: '--conf spark.executor.instances=2 --conf spark.executor.memory=4g' } }, configurationOverrides: { monitoringConfiguration: { s3MonitoringConfiguration: { logUri: `s3://${etlBucket.bucketName}/job-logs/` } } } }));
5. 注意事项
- 确保你的EMR Serverless应用处于
STARTED状态,否则任务无法执行。可以在CDK中添加启动应用的逻辑,或手动在控制台启动。 - 部署前验证执行角色的权限,避免因权限不足导致任务失败。
- 可通过CDK的
CfnOutput输出任务ID或ARN,方便后续在控制台或CLI中查看任务状态。
内容的提问来源于stack exchange,提问作者gospecomid12
相关产品推荐
相关产品推荐

