【问题标题】:Writing to multiple Kafka partitions from Spark从 Spark 写入多个 Kafka 分区
【发布时间】:2019-05-25 16:20:01
【问题描述】:

我有按此处指定的方式将批处理写入 Kafka 的 Spark 代码:

https://spark.apache.org/docs/2.4.0/structured-streaming-kafka-integration.html

代码如下所示:

  df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") 
   \
   .write \
   .format("kafka") \
   .option("kafka.bootstrap.servers", 
           "host1:port1,host2:port2") \
   .option("topic", "topic1") \
   .save()

但是数据只被写入 Kafka 分区 0。我怎样才能将它统一写入同一主题中的所有分区?

【问题讨论】:

  • 主题实际有多少个分区?
  • 主题有多少个分区? df 中有多少个不同的 keys?

标签: apache-spark apache-kafka


【解决方案1】:

Kafka 根据消息的密钥分发消息。因此,具有相同 key 的消息将被放置到同一个分区中。您的所有消息都可能具有相同的密钥。

【讨论】:

  • 这就是问题所在。感谢您的洞察力。
猜你喜欢
  • 2017-11-28
  • 2021-05-01
  • 2019-04-26
  • 2019-02-17
  • 2018-09-26
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-07-12
相关资源
最近更新 更多