【问题标题】:How to load all records that were already published from Kafka?如何加载已从 Kafka 发布的所有记录?
【发布时间】:2019-07-26 07:16:48
【问题描述】:

我有一个像这样设置的 pyspark 结构流式 python 应用程序

from pyspark.sql import SparkSession

spark = SparkSession\
    .builder\
    .appName("data streaming app")\
    .getOrCreate()


data_raw = spark.readStream\
    .format("kafka")\
    .option("kafka.bootstrap.servers", "kafkahost:9092")\
    .option("subscribe", "my_topic")\
    .load()

query = data_raw.writeStream\
    .outputMode("append")\
    .format("console")\
    .option("truncate", "false")\
    .trigger(processingTime="5 seconds")\
    .start()\
    .awaitTermination()

所有出现的都是这个

+---+-----+-----+---------+------+---------+-------------+
|key|value|topic|partition|offset|timestamp|timestampType|
+---+-----+-----+---------+------+---------+-------------+
+---+-----+-----+---------+------+---------+-------------+

19/03/04 22:00:50 INFO streaming.StreamExecution: Streaming query made progress: {
  "id" : "ab24bd30-6e2d-4c2a-92a2-ddad66906a5b",
  "runId" : "29592d76-892c-4b29-bcda-f4ef02aa1390",
  "name" : null,
  "timestamp" : "2019-03-04T22:00:49.389Z",
  "numInputRows" : 0,
  "processedRowsPerSecond" : 0.0,
  "durationMs" : {
    "addBatch" : 852,
    "getBatch" : 180,
    "getOffset" : 135,
    "queryPlanning" : 107,
    "triggerExecution" : 1321,
    "walCommit" : 27
  },
  "stateOperators" : [ ],
  "sources" : [ {
    "description" : "KafkaSource[Subscribe[my_topic]]",
    "startOffset" : null,
    "endOffset" : {
      "my_topic" : {
        "0" : 303
      }
    },
    "numInputRows" : 0,
    "processedRowsPerSecond" : 0.0
  } ],
  "sink" : {
    "description" : "org.apache.spark.sql.execution.streaming.ConsoleSink@74fad4a5"
  }
}

如您所见,my_topic 那里有 303 条消息,但我无法显示。其他信息包括我正在使用融合的 Kafka JDBC 连接器来查询 oracle 数据库并将行存储到 kafka 主题中。我有一个 avro 模式注册表设置。如果需要,我也会分享这些属性文件。

有人知道发生了什么吗?

【问题讨论】:

  • 你能检查一下broker的日志,确保没有错误吗?
  • @GiorgosMyrianthous,没必要,我想出了问题所在。这是我对流如何工作的有限理解......我想做的是显示自 kafka 主题首次启动以来的消息以及传入的消息。为此,我需要做的只是一个额外的选项,带有参数 startingOffsetsearliest

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


【解决方案1】:

作为流式应用程序,此 Spark 结构流式传输仅在消息发布后立即读取消息。为了测试目的,我想做的是阅读主题中的所有内容。为此,您只需在readStream 中添加一个额外选项,即option("startingOffsets", "earliest")

data_raw = spark.readStream\
    .format("kafka")\
    .option("kafka.bootstrap.servers", "kafkahost:9092")\
    .option("subscribe", "my_topic")\
    .option("startingOffsets", "earliest")
    .load()

【讨论】:

    猜你喜欢
    • 2019-11-04
    • 2020-02-10
    • 2021-09-03
    • 2021-12-06
    • 2019-07-04
    • 1970-01-01
    • 2010-10-17
    • 2022-01-07
    • 1970-01-01
    相关资源
    最近更新 更多