【发布时间】:2019-12-13 18:50:00
【问题描述】:
我使用 Apache Spark 2.4.1 和 kafka data source。
Dataset<Row> df = sparkSession
.readStream()
.format("kafka")
.option("kafka.bootstrap.servers", SERVERS)
.option("subscribe", TOPIC)
.option("startingOffsets", "latest")
.option("auto.offset.reset", "earliest")
.load();
我有两个接收器:原始数据存储在 hdfs 位置,经过几次转换,最终数据存储在 Cassandra 表中。 checkpointLocation 是一个 HDFS 目录。
启动流式查询时,它会给出以下警告:
2019-12-10 08:20:38,926 [任务 639 的执行程序任务启动工作人员] 警告 org.apache.spark.sql.kafka010.InternalKafkaConsumer - 一些数据 可能会丢失。从最早的偏移量恢复:470021 2019-12-10 08:20:38,926 [任务 639 的执行任务启动工作人员] WARN org.apache.spark.sql.kafka010.InternalKafkaConsumer - 当前 可用偏移范围为 AvailableOffsetRange(470021,470021)。抵消 62687 超出范围,将跳过 [62687, 62727) 中的记录 (组号: spark-kafka-source-1fba9e33-165f-42b4-a220-6697072f7172-1781964857-executor, 主题分区:INBOUND-19)。一些数据可能已经丢失,因为它们 在 Kafka 中不再可用;要么数据老化 Kafka 或主题可能在所有数据之前已被删除 主题已处理。如果您希望您的流式查询在此类上失败 在这种情况下,将源选项“failOnDataLoss”设置为“true”。
我还使用auto.offset.reset 作为latest 和startingOffsets 作为latest。
2019-12-11 08:33:37,496 [Executor task launch worker for task 1059] WARN org.apache.spark.sql.kafka010.KafkaDataConsumer - KafkaConsumer 缓存达到最大容量 64,删除 CacheKey(spark- kafka-source-93ee3689-79f9-42e8-b1ee-e856570205ae-1923743483-executor,_INBOUND-19)
这告诉我什么?如何摆脱警告(如果可能)?
【问题讨论】:
-
你能提供水槽的代码吗?顺便说一句,它们是两个不同的流式查询。这两个查询是否给出警告?
标签: apache-spark apache-kafka apache-spark-sql spark-structured-streaming