
点亮⭐️
https://github.com/apache/
点击蓝字 关注我们
01 项目概述:当DolphinScheduler
遇上AWS数据湖仓
在数据驱动的业务决策成为常态的今天,一个高效、灵活且成本可控的数据处理流水线是企业数据中台的核心竞争力。我接触过不少团队,他们早期往往在本地数据中心搭建Hadoop集群,自己维护一套调度系统,但随着数据量激增和业务复杂度提升,这套体系的瓶颈日益凸显:资源扩容周期长、运维成本高、任务调度不够灵活,最关键的是,很难在性能和成本之间找到一个优雅的平衡点。
近年来,云原生数据架构的成熟,特别是像AWS这样提供从数据湖(S3)到数据仓库(Redshift)、再到大数据处理引擎(EMR)的全栈服务,为这个问题提供了新的解题思路。然而,仅仅把组件搬到云上还不够,如何将它们像乐高积木一样无缝拼接、统一调度,才是真正释放云平台潜力的关键。这正是 Apache DolphinScheduler 作为一款开源的分布式可视化工作流任务调度平台,能够大显身手的地方。它就像一个“数据流水线的总指挥”,能够协调AWS EMR进行大规模数据处理,并调度Redshift完成高性能分析查询。
本文将基于一个真实的客户迁移与优化案例,深入探讨如何将DolphinScheduler深度集成到AWS的智能数据湖仓架构中。我们会重点拆解两个核心场景:一是如何利用DolphinScheduler对AWS EMR(包括传统的集群模式和Serverless无服务器模式)进行混合调度与成本优化;二是如何通过DolphinScheduler高效、安全地调度AWS Redshift任务,并解决其固有的并发限制挑战。无论你是正在规划上云的数据平台工程师,还是寻求现有云上数据流水线效能提升的架构师,相信这些从一线实战中总结出的模式、代码和避坑经验,都能给你带来直接的参考价值。
02 架构与选型:为什么是
DolphinScheduler + AWS组合?
在深入实操之前,我们有必要先厘清整个技术栈的选型逻辑。为什么是AWS?为什么是DolphinScheduler?这个组合解决了哪些痛点?理解了这些“为什么”,后续的具体实施才会更有方向。
AWS提出的智能数据湖仓(Intelligent Data Lakehouse)并非一个单一产品,而是一个以Amazon S3数据湖为中心的最佳实践架构。它的核心思想是打破数据孤岛,让数据在存储、处理、分析、机器学习各环节间自由、安全地流动。
在这个架构中,S3作为统一的、无限扩展的数据存储层,存放所有原始和加工后的数据。围绕S3,一系列托管服务各司其职:
数据摄取
大数据处理:Amazon EMR 在这里扮演核心角色,它托管了Spark、Hive、Flink等开源框架,用于执行数据清洗、转换和复杂计算。
数据仓库与分析:Amazon Redshift 作为高性能云数据仓库,负责对处理后的数据进行快速、复杂的查询分析,支撑BI报表和即席查询。
机器学习与AI:Amazon SageMaker等服务可以直接读取S3或Redshift中的数据,进行模型训练和推理。
这个架构的“智能”之处在于,所有服务都是深度集成的。例如,EMR和Redshift都可以直接读取S3中的数据,无需移动;Glue Data Catalog可以作为统一的元数据中心,被EMR、Redshift和Athena共享。而我们的挑战在于,如何用一个统一的调度器,将散落在这些服务上的任务有序、可靠、高效地串联起来。
Apache DolphinScheduler是一个分布式、可视化的工作流调度系统。与传统的Crontab或简单的脚本调度相比,它的优势非常明显:
可视化编排:通过拖拽方式定义复杂的DAG(有向无环图)工作流,依赖关系一目了然,降低了数据开发的门槛。
高可靠与高可用:采用去中心化的Master-Worker架构,支持水平扩展,任一节点故障不会导致服务不可用,任务支持失败重试、告警等机制。
多租户与资源隔离:支持按项目、用户进行资源隔离和权限控制,适合团队协作。
丰富的任务类型:原生支持Shell、Python、SQL、Spark、Flink等多种任务类型,并且具备强大的 插件扩展能力 ,这正是我们能将其与AWS服务深度集成的基石。
在AWS数据湖仓架构中,DolphinScheduler的定位就是 “编排层” 。它不替代EMR的计算能力,也不替代Redshift的查询能力,而是站在更高的维度,决定在什么时间、用什么参数、将哪个任务(Spark作业、Redshift SQL)提交到哪个计算资源(EMR集群、EMR Serverless、Redshift集群)上执行,并监控其全生命周期。
客户从本地Hadoop迁移到AWS EMR后,最初全部使用“EMR on EC2”(即托管在EC2虚拟机上的集群)模式。但他们很快发现,任务负载存在明显的波峰波谷和大小任务差异,导致集群资源利用率不均,成本有优化空间。
我们面临的典型任务负载分为三类:
小型任务:数量多,单个执行时间短(20-30分钟),全天候分散触发。这类任务启动一个完整的EMR集群(通常至少3-4个节点)来运行,好比“用高射炮打蚊子”,资源浪费严重。
大型任务:每日在固定的7-8小时窗口内集中运行,计算密集。适合用专用集群处理,但集群在非窗口期闲置会产生费用。
超大型任务:执行周期长(数天),每月仅运行2-3次。为它们长期维护一个集群极不经济。
基于此,我们制定的混合调度策略是:
小型 & 超大型任务 → EMR Serverless:EMR Serverless是AWS推出的无服务器模式,你无需预置或管理集群,只需提交作业,按实际消耗的vCPU和内存付费。它完美契合了 sporadic(零星)和 bursty(突发)的工作负载。小型任务避免了集群空转成本,超大型任务则避免了为偶发任务预留巨额资源。
大型任务 → EMR on EC2(定时启停):对于每日固定窗口的密集型任务,我们继续使用EMR on EC2,但通过DolphinScheduler在任务开始前自动启动集群,任务结束后自动终止集群。结合使用Spot实例(竞价实例)可以进一步压缩60%-70%的计算成本。
这个策略的核心在于,DolphinScheduler需要具备智能路由的能力,能根据任务属性(如标签、资源需求)自动决定将其提交到Serverless还是集群模式。这要求我们对DolphinScheduler进行定制化开发,封装统一的提交接口。
03 核心集成实战:封装统一API调度EMR
理论很美好,但落地时EMR on EC2和EMR Serverless两套服务在API、交互方式上的差异,是第一个需要跨过的坎。我们的目标是让数据开发人员在DolphinScheduler上编写任务时,无需关心底层是哪种EMR,实现透明化调度。
在集成过程中,我们遇到了几个关键的技术差异点:
任务提交模式:
日志查看方式:
API接口差异:
SQL支持度:
如果让用户在DolphinScheduler任务里写两套代码,无疑增加了复杂度和维护成本。我们的解决方案是:封装一个统一的Python SDK。
我们开发了一个名为 emr_common 的Python库,核心是提供一个统一的 Session 类。这个类根据传入的 job_type 参数,在内部实例化不同的子会话对象( EMRSession 或 EMRServerlessSession ),但对外暴露完全一致的接口。
# emr_common/session.pyimport boto3from abc import ABC, abstractmethodclass BaseEMRSession(ABC): """EMR会话抽象基类""" @abstractmethod def submit_sql(self, job_name: str, sql: str, **kwargs): pass @abstractmethod def submit_file(self, job_name: str, file_path: str, **kwargs): pass @abstractmethod def get_status(self, job_id: str) -> str: pass @abstractmethod def get_logs(self, job_id: str) -> str: passclass EMRSession(BaseEMRSession): """EMR on EC2 会话实现""" def __init__(self, cluster_id: str = None): self.client = boto3.client('emr') # 如果未指定集群,则查找一个正在运行的集群 self.cluster_id = cluster_id or self._find_active_cluster() def submit_sql(self, job_name: str, sql: str, **kwargs): # 构造EMR Step参数 step_args = [ 'spark-sql', '-e', sql ] response = self.client.add_job_flow_steps( JobFlowId=self.cluster_id, Steps=[{ 'Name': job_name, 'ActionOnFailure': 'CONTINUE', 'HadoopJarStep': { 'Jar': 'command-runner.jar', 'Args': step_args } }] ) return response['StepIds'][0] # ... 其他方法实现(submit_file, get_status, get_logs)class EMRServerlessSession(BaseEMRSession): """EMR Serverless 会话实现""" def __init__(self, application_id: str = None): self.client = boto3.client('emr-serverless') self.application_id = application_id or self._find_active_application() def submit_sql(self, job_name: str, sql: str, **kwargs): # EMR Serverless 不支持直接提交SQL字符串,需要封装成脚本 # 我们将SQL写入一个临时PySpark脚本文件,然后提交该文件 import tempfile with tempfile.NamedTemporaryFile(mode='w', suffix='.py', delete=False) as f: f.write(f"""from pyspark.sql import SparkSessionspark = SparkSession.builder.appName("{job_name}").getOrCreate()spark.sql(\"\"\"{sql}\"\"\")""") script_path = f.name # 上传脚本到S3(此处省略上传代码) s3_script_uri = self._upload_to_s3(script_path) return self.submit_file(job_name, s3_script_uri) def submit_file(self, job_name: str, file_path: str, **kwargs): # 构造EMR Serverless作业请求 response = self.client.start_job_run( applicationId=self.application_id, executionRoleArn='arn:aws:iam::xxx:role/EMRServerlessRole', jobDriver={ 'sparkSubmit': { 'entryPoint': file_path, # 可以在这里传递Spark配置 'sparkSubmitParameters': '--conf spark.executor.memory=4g' } }, configurationOverrides={...}, name=job_name ) return response['jobRunId'] # ... 其他方法实现(get_status, get_logs)需处理异步轮询和S3日志获取def Session(job_type: int = 0, **kwargs): """工厂函数,返回统一的会话对象""" if job_type == 0: return EMRSession(**kwargs) elif job_type == 1: return EMRServerlessSession(**kwargs) else: raise ValueError(f"Unsupported job_type: {job_type}")关键设计要点 :
同步化异步接口:在 EMRServerlessSession.get_status() 方法内部,我们实现了一个轮询机制,阻塞调用直到作业完成(或失败),从而对上层提供同步的体验。
日志统一拉取:get_logs() 方法内部判断作业类型,如果是Serverless则从S3下载并拼接日志文件,如果是EMR on EC2则从CloudWatch或直接SSH到主节点获取,最终返回格式统一的日志文本。
默认值处理:例如,当用户未指定 cluster_id 或 application_id 时,SDK会自动查找当前账户/区域下处于 WAITING 或 STARTED 状态的集群/应用,降低了用户的配置负担。
Spark参数统一:我们设计了一个统一的配置字典,可以传递 driver_memory , executor_cores 等参数,SDK内部会将其翻译成对应服务所需的参数格式。
封装好SDK后,在DolphinScheduler中使用就变得异常简单。我们主要使用 Python Operator 来调用这个SDK。
在DolphinScheduler中创建Python任务节点:
# 示例:一个可根据参数动态选择EMR类型的任务from emr_common import Sessiondef main(**kwargs): # 从DolphinScheduler系统参数或上游节点获取任务类型 # 这里假设通过自定义参数 `emr_job_type` 传递,0代表EMR on EC2, 1代表EMR Serverless job_type = int(kwargs.get('emr_job_type', 0)) task_name = kwargs.get('task_name', 'default_task') # 创建统一会话 session = Session(job_type=job_type) # 示例1:提交一个SQL任务 sql = """ SELECT date, count(*) as cnt FROM your_table WHERE date = '${bizdate}' -- DolphinScheduler内置参数,替换为业务日期 GROUP BY date """ job_id = session.submit_sql(job_name=f"{task_name}_sql", sql=sql) # 等待任务完成并获取状态 status = session.wait_for_completion(job_id, interval=30) # 每30秒检查一次 print(f"Job {job_id} finished with status: {status}") # 获取并打印日志(可用于任务失败时排查) logs = session.get_logs(job_id) # 可以将关键日志写入DolphinScheduler任务日志 # ... # 根据状态决定任务成功或失败 if status == 'SUCCESS': return True else: raise Exception(f"Job failed with status: {status}")if __name__ == "__main__": # DolphinScheduler会将参数通过字典传入 import sys import json # 解析参数,这里是一个简化示例 params = json.loads(sys.argv[1]) if len(sys.argv) > 1 else {} success = main(**params) sys.exit(0 if success else 1)工作流参数化设计: 我们可以在DolphinScheduler的工作流定义中,设置全局参数或节点参数。例如,定义一个名为 emr_engine 的全局参数,默认值为 serverless 。在Python任务中,通过 ${emr_engine} 引用。我们甚至可以写一个简单的映射函数,将 serverless 映射为 job_type=1 。这样,只需修改工作流全局参数,就能一键切换整个工作流所有任务的执行引擎,极大地提升了灵活性和测试效率。
关于元数据统一: 为了让EMR on EC2和EMR Serverless能够无缝读写同一份数据并使用同一套元数据,我们强烈推荐使用 AWS Glue Data Catalog 作为统一的元数据存储。在创建EMR集群或Serverless应用时,都将其配置为使用Glue Catalog。这样,无论任务在哪种引擎上运行,它们看到的库、表结构都是一致的,Hive Metastore的维护难题也迎刃而解。
04 Redshift任务调度与并发控制实战
集成完EMR,我们来看数据仓库层——Amazon Redshift。Redshift是一款成熟的PB级云数据仓库,擅长复杂查询和快速分析。从DolphinScheduler 3.x版本开始,其内置的 数据源中心 已经支持直接添加Redshift数据源,这为我们通过SQL任务直接操作Redshift铺平了道路。
这是最直接、最常用的方式。在DolphinScheduler的UI上配置好Redshift数据源(需要JDBC驱动和网络连通性)后,就可以创建SQL任务节点。
配置步骤与要点:
-- 示例:创建一个每日分区表并插入数据CREATE TABLE IF NOT EXISTS dws_user_daily ( biz_date DATE, user_id BIGINT, activity_count INT) DISTKEY(user_id) SORTKEY(biz_date);DELETE FROM dws_user_daily WHERE biz_date = '${bizdate}';INSERT INTO dws_user_dailySELECT '${bizdate}'::DATE as biz_date, user_id, COUNT(*) as activity_countFROM ods_user_logWHERE event_date = '${bizdate}'GROUP BY user_id;注意 :Redshift对事务的支持与OLTP数据库不同, BEGIN; ... COMMIT; 在DDL和大量DML操作中行为有差异。建议在调度中将DDL(CREATE, ALTER, DROP)和DML(INSERT, UPDATE, DELETE)分开成不同的任务节点,并设置好依赖关系,避免锁表和意外回滚。
Redshift基于MPP架构,虽然查询速度快,但一个集群的并发查询槽位(concurrency slots)是有限的,默认通常不超过50。当DolphinScheduler调度大量任务同时涌向Redshift时,很容易导致并发超限,任务排队甚至失败。
我们实践过两种有效的控制策略:
这是Redshift提供的一项付费功能。当主集群的并发槽位用尽时,Redshift会自动启动临时的扩展集群来处理增加的负载,最高可扩展至10倍。优势是“无感”扩容,对代码和调度器零改造。
这是更经济、更可控的方案。DolphinScheduler支持创建“任务组”并设置组的最大并发数。
两种策略对比与选型建议:

我们的最佳实践是两者结合:为日常调度任务设置DolphinScheduler任务组进行基线控制;同时为Redshift集群启用1-2个并发扩展集群作为缓冲,以应对临时增加的即席查询或某个调度任务异常复杂导致的长时间占用。
除了SQL Operator, Shell Operator 也是一个强大的工具,特别是当你的SQL脚本已经文件化,并希望通过Git进行版本控制时。
典型使用模式 :
# 假设Redshift的`psql`命令行工具已安装在Worker节点上,或使用AWS CLI执行export PGPASSWORD=${redshift_password}psql -h ${redshift_host} -p ${redshift_port} -U ${redshift_user} -d ${redshift_db} \ -f s3://your-bucket/scripts/transform_sales_daily.sql \ -v bizdate=\'${bizdate}\'-f参数指定从S3读取SQL文件(需确保Redshift集群有访问该S3桶的权限,通常通过IAM角色)。-v
与DolphinScheduler资源中心集成:更优雅的方式是利用DolphinScheduler的 资源中心 功能。你可以将S3桶挂载为DolphinScheduler的一个存储资源(需要实现或使用支持S3的文件存储插件)。这样,SQL脚本文件可以直接在DolphinScheduler的UI上进行管理、编辑和版本查看。在Shell任务中,直接引用资源中心内的文件路径即可,DolphinScheduler会自动处理文件的拉取。
这种模式将调度(DolphinScheduler)、代码管理(Git)、持续集成(Jenkins)和对象存储(S3)紧密结合起来,实现了数据开发流程的DevOps化,保证了脚本的版本一致性和可追溯性。
05 运维、监控与成本优化经验谈
将调度系统与云服务深度集成后,运维和监控视角也需要从单个服务扩展到整个链路。同时,云上资源“按需付费”的特性,使得成本优化成为一个持续的过程,而不仅仅是一次性的迁移动作。
一个任务失败,可能是DolphinScheduler自身问题,可能是网络问题,可能是EMR或Redshift资源不足,也可能是脚本逻辑错误。我们需要建立端到端的监控。
云上成本优化是一个精细活。以下是我们从客户案例中总结出的针对本架构的优化点:
针对EMR的优化:
针对Redshift的优化:
在实际运维中,以下问题较为常见:

06 总结与展望
回顾整个集成实践,其核心价值在于通过DolphinScheduler这一层“抽象”,将AWS上异构的计算服务(EMR on EC2, EMR Serverless, Redshift)整合成了一个逻辑统一的 数据计算平台 。数据开发人员只需关注业务逻辑(写SQL或PySpark),而无需纠结于任务该在哪里运行、集群如何管理。运维人员则获得了全局的调度视图、统一的监控和成本控制抓手。
这次深度集成的成功,离不开两个关键设计:一是面向接口的SDK封装,它抹平了底层服务的差异;二是充分利用了DolphinScheduler的参数化、插件化和资源控制能力,使得调度策略可以灵活定制。
对于未来,随着DolphinScheduler社区的不断发展,我期待能在两个方面看到更多进展,这也会让云上数据流水线更加智能和高效:
基于SQL语法树的数据血缘解析:
工作流编排中引入AI智能体:
云原生的道路没有终点,工具链的深度集成与智能化是提升数据团队产能的关键。希望本文分享的具体方案、代码片段和踩坑经验,能为你构建或优化自己的数据调度平台提供切实可行的参考。
原文链接:https://blog.csdn.net/weixin_30847865/article/details/95634249
END

用户案例

迁移实战

最新发版消息

加入社区
关注社区的方式有很多:
同样地,参与Apache DolphinScheduler 有非常多的参与贡献的方式,主要分为代码方式和非代码方式两种。
非代码方式包括:
完善文档、翻译文档;翻译技术性、实践性文章;投稿实践性、原理性文章;成为布道师;社区管理、答疑;会议分享;测试反馈;用户反馈等。
代码方式包括:
查找Bug;编写修复代码;开发新功能;提交代码贡献;参与代码审查等。


你的好友小海豚拍了拍你
并请你帮她点一下“分享”
