【发布时间】:2020-10-17 21:43:08
【问题描述】:
我正在尝试运行多个 Spark Structured Streaming 作业(在 EMR 上),这些作业从 Kafka 主题读取并写入 S3 中的不同路径(每个都在各自的作业中执行)。我已将集群配置为使用CapacityScheduler。这是我尝试运行的代码的 sn-p:
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", <BOOTSTRAP_SERVERS>) \
.option("subscribePattern", "<MY_TOPIC>") \
.load() \
.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
output = df \
.writeStream \
.format("json") \
.outputMode("update") \
.option("checkpointLocation", "s3://<CHECKPOINT_LOCATION>") \
.option("path", "s3://<SINK>") \
.start() \
.awaitTermination()
我尝试并行运行两个作业:
spark-submit --queue <QUEUE_1> --deploy-mode cluster --master yarn <STREAM_1_SCRIPT>.py
spark-submit --queue <QUEUE_2> --deploy-mode cluster --master yarn <STREAM_2_SCRIPT>.py
在执行期间,我注意到第二个作业没有写入 S3(即使第一个作业是)。我还注意到第二个作业通过 Spark UI 的利用率出现了巨大的峰值。
停止第一个作业后,数据显示在 S3 中的第二个作业。 不可能运行两个并行写入接收器(特别是在 S3 上)的独立 Spark 结构化流式处理作业吗?写操作是否会导致某种阻塞?
【问题讨论】:
-
.py 和 .py 有什么区别? -
两者是否使用相同的检查点位置??
-
两者都写入同一个 s3 位置吗?
-
.py 和 .py 包含上面代码的 sn-p。唯一的区别在于实际的主题名称、检查点位置和 S3 路径(接收器)。他们没有使用相同的检查点位置。 -
我想拥有不同的工作,这样我就可以将主题的摄取彼此分离。如果某个主题发生了某些事情(例如错误),我不想不得不停止对另一个主题的摄取(因为流在一个单独的流式作业中会发生这种情况)。
标签: amazon-web-services apache-spark amazon-emr spark-structured-streaming