Advertisement

Streamling-Kafka-Spark-流处理-Kafka集成

阅读量:

近期在使用spark-streaming-kafka过程中遇到的一些问题

问题1

存在部分消息的尺寸超过当前设定的fetch size(1048576),因此这些消息将无法被正常获取。解决方式包括提升fetch size的数值,或者降低broker允许的最大消息尺寸。

解决方案

通过调整ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG(即fetch.message.max.bytes)参数,将其设置为一个更高的数值以应对该问题。

问题2

问题描述:当kafka中堆积了大量尚未被消费的消息时

在启动spark-streaming程序以消费kafka中的数据时,若kafka中存在大量未处理的消息,则首个batch需要处理的数据量会非常庞大,这可能会超出spark-submit所配置资源的承载能力,从而引发程序崩溃。为了解决这一情况,在执行spark-submit命令时额外添加了一个配置项:--conf spark.streaming.kafka.maxRatePerPartition=10000。该参数用于限制每个partition每秒钟最多可消费的消息数量,以此将原本集中于首个batch的大批量消息分散至多个batch中进行处理。为了更高效地处理延迟的消息,可以适当增加计算资源并提高该参数的取值范围。

spark exec

全部评论 (0)

还没有任何评论哟~