【问题标题】:spark streaming kafka : Unknown error fetching data for topic-partitionspark streaming kafka:为主题分区获取数据时出现未知错误
【发布时间】:2018-10-29 00:23:33
【问题描述】:

我正在尝试使用结构化流 API 和 Spark 中的 Kafka 集成从 Spark 集群中读取 Kafka 主题

val sparkSession = SparkSession.builder()
  .master("local[*]")
  .appName("some-app")
  .getOrCreate()

Kafka 流创建

import sparkSession.implicits._

val dataFrame = sparkSession
  .readStream
  .format("kafka")
  .option("subscribepattern", "preprod-*")
  .option("kafka.bootstrap.servers", "<brokerUrl>:9094")
  .option("kafka.ssl.protocol", "TLS")
  .option("kafka.security.protocol", "SSL")
  .option("kafka.ssl.key.password", secretPassword)
  .option("kafka.ssl.keystore.location", "/tmp/xyz.jks")
  .option("kafka.ssl.keystore.password", secretPassword)
  .option("kafka.ssl.truststore.location", "/abc.jks")
  .option("kafka.ssl.truststore.password", secretPassword)
  .load()
  .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
  .as[(String, String)]
  .writeStream
  .format("console")
  .start()
  .awaitTermination()

使用命令运行它

/usr/local/spark/bin/spark-submit 
--packages "org.apache.spark:spark-streaming-kafka-0-10_2.11:2.3.1,org.apache.spark:spark-sql-kafka-0-10_2.11:2.3.1"
myjar.jar

得到以下错误

    2018-09-28 07:29:23 INFO  AbstractCoordinator:505 - Discovered coordinator brokerUrl.com:32400 (id: 2147483647 rack: null) for group spark-kafka-source-c72dcb79-f3bc-4dfd-86a5-9d14be48fa04-1188588017-executor.
2018-09-28 07:29:23 INFO  AbstractCoordinator:505 - Discovered coordinator brokerUrl.com:32400 (id: 2147483647 rack: null) for group spark-kafka-source-c72dcb79-f3bc-4dfd-86a5-9d14be48fa04-1188588017-executor.
2018-09-28 07:29:23 INFO  AbstractCoordinator:505 - Discovered coordinator brokerUrl.com:32400 (id: 2147483647 rack: null) for group spark-kafka-source-c72dcb79-f3bc-4dfd-86a5-9d14be48fa04-1188588017-executor.
2018-09-28 07:29:23 INFO  AbstractCoordinator:505 - Discovered coordinator brokerUrl.com:32400 (id: 2147483647 rack: null) for group spark-kafka-source-c72dcb79-f3bc-4dfd-86a5-9d14be48fa04-1188588017-executor.
2018-09-28 07:29:47 WARN  Fetcher:594 - Unknown error fetching data for topic-partition preprod-sanity-test-5
2018-09-28 07:30:25 WARN  Fetcher:594 - Unknown error fetching data for topic-partition preprod-sanity-test-7
2018-09-28 07:30:27 WARN  Fetcher:594 - Unknown error fetching data for topic-partition preprod-sanity-test-7
2018-09-28 07:30:27 WARN  Fetcher:594 - Unknown error fetching data for topic-partition preprod-sanity-test-5
2018-09-28 07:30:50 WARN  Fetcher:594 - Unknown error fetching data for topic-partition preprod-sanity-test-8
2018-09-28 07:30:50 WARN  Fetcher:594 - Unknown error fetching data for topic-partition preprod-sanity-test-4
2018-09-28 07:30:50 WARN  Fetcher:594 - Unknown error fetching data for topic-partition preprod-sanity-test-7
2018-09-28 07:30:50 WARN  Fetcher:594 - Unknown error fetching data for topic-partition preprod-sanity-test-8
2018-09-28 07:30:50 WARN  Fetcher:594 - Unknown error fetching data for topic-partition preprod-sanity-test-4
2018-09-28 07:30:50 WARN  Fetcher:594 - Unknown error fetching data for topic-partition preprod-sanity-test-5
.....
....
so on

【问题讨论】:

    标签: apache-spark apache-kafka apache-spark-sql spark-structured-streaming


    【解决方案1】:

    您的 Kafka 代理版本是什么?您是如何生成这些消息的?

    如果这些消息有标头 (https://issues.apache.org/jira/browse/KAFKA-4208),您将需要使用 Kafka 0.11+ 来使用它们,因为旧的 Kafka 客户端无法读取此类消息。如果是这样,您可以使用以下命令:

    /usr/local/spark/bin/spark-submit --packages "org.apache.kafka:kafka-clients:0.11.0.3,org.apache.spark:spark-sql-kafka-0-10_2.11:2.3.1"
    myjar.jar
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-06-24
      • 2016-06-09
      • 2018-09-26
      • 2017-09-30
      • 2019-07-12
      • 1970-01-01
      • 2023-03-25
      • 2018-09-06
      相关资源
      最近更新 更多