如何在Dataproc工作流执行完成后向Pub/Sub主题发送通知?
实现Dataproc工作流结束后向Pub/Sub发送消息的方案
方案一:在工作流末尾添加自定义消息发送作业
这是最直接的方式,把发送Pub/Sub消息作为工作流的最后一步,确保所有Hadoop、Pig、Spark作业完成后才执行。
- 编写一个简单的消息发送脚本(以Python为例):
from google.cloud import pubsub_v1 import sys def send_message(project_id, topic_id, message): publisher = pubsub_v1.PublisherClient() topic_path = publisher.topic_path(project_id, topic_id) data = message.encode("utf-8") future = publisher.publish(topic_path, data) print(f"Message published: {future.result()}") if __name__ == "__main__": if len(sys.argv) != 4: print("Usage: python send_pubsub.py <project-id> <topic-id> <message>") sys.exit(1) project_id = sys.argv[1] topic_id = sys.argv[2] message = sys.argv[3] send_message(project_id, topic_id, message)
将脚本上传到Cloud Storage(示例路径:
gs://your-bucket/scripts/send_pubsub.py)在Dataproc工作流中添加最后一步作业:
- 作业类型选择
PYSPARK或CUSTOM - 配置主文件路径为上述Cloud Storage中的脚本地址
- 传入参数:
<你的项目ID> <Pub/Sub主题ID> "Dataproc工作流已完成" - 设置该作业依赖前面所有Hadoop、Pig、Spark作业,确保仅当前面作业全部成功后才执行
- 作业类型选择
方案二:通过Cloud Function监听Dataproc作业状态
利用Cloud Eventarc监听Dataproc作业生命周期事件,当作业进入DONE状态时触发Cloud Function发送消息。
- 创建Cloud Function:
- 运行环境选择Python 3.x
- 触发器类型选
Cloud Event,事件源选Cloud Dataproc,事件类型选google.cloud.dataproc.job.v1.completed - 编写函数代码:
from google.cloud import pubsub_v1 def send_pubsub_on_job_done(event, context): job_data = event.get("data", {}) job_status = job_data.get("status", {}).get("state") # 仅处理成功完成的作业(可根据需求调整逻辑) if job_status == "DONE" and job_data.get("status", {}).get("stateDetails") == "SUCCEEDED": project_id = "你的项目ID" topic_id = "你的Pub/Sub主题ID" message = f"Dataproc作业 {job_data.get('reference', {}).get('jobId')} 已完成" publisher = pubsub_v1.PublisherClient() topic_path = publisher.topic_path(project_id, topic_id) data = message.encode("utf-8") publisher.publish(topic_path, data)
- 配置权限:
- 为Cloud Function分配
Pub/Sub Publisher角色,使其具备向目标主题发消息的权限 - 确保Eventarc拥有触发该函数的权限
- 为Cloud Function分配
方案三:利用Dataproc工作流模板的依赖链
如果使用Dataproc Workflow Templates管理工作流,可通过定义作业依赖将消息发送设置为最终节点:
- 创建Workflow Template时,依次添加Hadoop、Pig、Spark作业,最后添加自定义发送消息作业
- 为发送消息作业设置
depends_on参数,指定前面所有作业的ID,确保工作流按顺序执行,最后一步发送消息
注意事项
- 确保Dataproc集群服务账号拥有
Pub/Sub Publisher角色,具备访问Pub/Sub的权限 - 若需处理作业失败场景,可在消息中包含作业状态信息,或调整逻辑仅在作业成功完成时发送消息
内容的提问来源于stack exchange,提问作者Tarique
相关产品推荐
相关产品推荐

