AWS Glue作业脚本能否在同一作业中处理多个流数据源?
解决AWS Glue合并多流作业仅单任务生效的问题
核心问题
Spark Structured Streaming 默认以单线程方式执行流查询,即便设置了spark.scheduler.mode=FAIR,若未正确配置调度池或未异步启动多流查询,代码会阻塞在第一个流的awaitTermination()调用上,导致后续流任务无法启动。
可行解决方案
1. 配置公平调度池并为每个流分配独立池
首先在SparkSession初始化时开启公平调度,并指定调度配置文件:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("MultiStreamGlueJob") \ .config("spark.scheduler.mode", "FAIR") \ .config("spark.scheduler.allocation.file", "/opt/spark/conf/fairscheduler.xml") \ .getOrCreate()
创建fairscheduler.xml文件(需上传至Glue作业可访问的路径,如S3并挂载),内容示例:
<allocations> <pool name="kafka-stream-1"> <schedulingMode>FAIR</schedulingMode> <weight>1</weight> </pool> <pool name="kafka-stream-2"> <schedulingMode>FAIR</schedulingMode> <weight>1</weight> </pool> </allocations>
然后为每个流查询指定对应的调度池:
# 处理第一个Kafka流 spark.sparkContext.setLocalProperty("spark.scheduler.pool", "kafka-stream-1") query1 = raw_data_frame_one.writeStream \ .foreachBatch(batch_function_one) \ .start() # 处理第二个Kafka流 spark.sparkContext.setLocalProperty("spark.scheduler.pool", "kafka-stream-2") query2 = raw_data_frame_two.writeStream \ .foreachBatch(batch_function_two) \ .start()
2. 异步等待所有流查询终止
不要仅调用单个查询的awaitTermination(),而是通过线程或Spark内置方法等待所有流:
方式一:使用多线程等待
from threading import Thread def wait_for_query(query): query.awaitTermination() # 启动独立线程等待每个流 Thread(target=wait_for_query, args=(query1,)).start() Thread(target=wait_for_query, args=(query2,)).start()
方式二:使用Spark 3.0+的awaitAnyTermination()
# 等待任意一个流终止,或通过循环等待所有 spark.streams.awaitAnyTermination()
Glue环境额外注意事项
- 确保作业Worker的数量和规格足够支撑多流并行,避免资源竞争导致某流无法执行。
- 每个Kafka流使用独立的消费组ID,避免消费offset冲突。
- Hudi表配置需完全独立(表名、存储路径、主键等),防止写入冲突。
内容的提问来源于stack exchange,提问作者Colin
相关产品推荐
相关产品推荐

