【问题标题】:Is it possible to run multiple Spark Structured Streaming jobs that write to S3 in parallel?是否可以运行多个并行写入 S3 的 Spark Structured Streaming 作业?
【发布时间】: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


【解决方案1】:

是的,你可以! 那不是记录了多个来源的东西,而是唯一需要它在多个作业的线程之间共享火花上下文的东西。 我在这篇文章https://cm.engineering/multiple-spark-streaming-jobs-in-a-single-emr-cluster-ca86c28d1411 之后制作了一个多火花结构化流式传输管道,您可以给我发送电子邮件或与我交谈。

谢谢!

【讨论】:

猜你喜欢
  • 2023-03-13
  • 2016-09-21
  • 2016-08-31
  • 2020-03-19
  • 2017-06-19
  • 2019-01-19
  • 1970-01-01
  • 2021-12-07
  • 2017-04-22
相关资源
最近更新 更多