You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 22:05:18