【发布时间】:2021-02-28 09:04:56
【问题描述】:
我们使用的是带有 5 个节点的小型 spark 集群,所有这 5 个节点都与 Kafka 代理连接。
我们计划通过添加更多节点来扩展集群,这可能需要配置额外的节点以与 Kafka 集群连接。我们正在评估集成的最佳实践
- 实际如何集成以使集成尽可能简单
- 是否需要所有工作节点都与 经纪人,在那种情况下,它可能无法扩展?
【问题讨论】:
我们使用的是带有 5 个节点的小型 spark 集群,所有这 5 个节点都与 Kafka 代理连接。
我们计划通过添加更多节点来扩展集群,这可能需要配置额外的节点以与 Kafka 集群连接。我们正在评估集成的最佳实践
【问题讨论】:
我建议使用 kafka 集成查看 spark 的文档
https://spark.apache.org/docs/latest/structured-streaming-kafka-integration.html
“如何实际集成以使集成尽可能简单”:
我不确定你的意思是什么——但基本上当你连接到 kafka 时,你应该提供引导服务器:引导服务器是主机/端口对的列表,用于建立与 Kafka 集群的初始连接。 这些服务器仅用于初始连接以发现完整的集群成员。所以kafka集群的节点数不会改变你集成的方式
“是否需要所有工作节点都与代理连接,在这种情况下,它可能无法扩展?” :
spark 集成的工作方式如下(某种程度上):
旁注:您可以使用配置来进一步打破将从单个 kafka 分区读取的 spark 分区的数量 - 它称为 minPartitions 并且来自 spark 2.4.7
最后一点:使用 kafka 的 spark 流式传输是一个非常常用且广为人知的用例,并且在非常大的数据生态系统中使用,作为第一个直观的想法,我认为它是可扩展的
【讨论】:
在翻书的过程中发现了以下短语,https://learning.oreilly.com/library/view/stream-processing-with/9781491944233/ch19.html
尤其是短语The driver does not send data to the executors; instead, it simply sends a few offsets they use to directly consume data. - 似乎_所有执行程序(工作节点)都必须与kafka连接,因为任务很可能在任何执行程序上运行
数据传递的要点是 Spark 驱动程序查询偏移量和 确定来自 Apache Kafka 的每个批处理间隔的偏移范围。 收到这些偏移量后,驱动程序通过 为每个分区启动一个任务,从而实现 1:1 并行度 工作中的 Kafka 分区和 Spark 分区之间。每个 任务使用其特定的偏移范围检索数据。
驱动程序确实 不向执行者发送数据;相反,它只是发送一些偏移量 他们用来直接消费数据。因此,并行性 来自 Apache Kafka 的数据摄取量比传统的要好得多 接收器模型,其中每个流由单个机器使用。
【讨论】: