【问题标题】:Spark structured streaming truncates Kafka value string to 4095Spark 结构化流将 Kafka 值字符串截断为 4095
【发布时间】:2021-07-02 06:30:40
【问题描述】:

以下代码

builder = SparkSession.builder\
   .appName("PythonTest11")
spark = builder.getOrCreate()

#spark.conf.set("spark.sql.debug.maxToStringFields", 10000)

# Subscribe to 1 topic
df = spark \
  .readStream \
  .format("kafka") \
  .option("kafka.bootstrap.servers", config["kafka"]["bootstrap.servers"]) \
  .option("subscribe", dataFlowTopic) \
  .load()

df = df \
    .selectExpr("LENGTH(value)")
#.selectExpr("CAST(value as string)") \


df.printSchema()

# Start running the query that prints the running counts to the console
query = df \
    .writeStream \
    .outputMode('append') \
    .format('console') \
    .start()

query.awaitTermination()

打印

+-------------+
|length(value)|
+-------------+
|         4095|
+-------------+

对于任何大消息,即截断传入的字符串。

如何解决这个问题?

【问题讨论】:

  • 您选择的是字节数组的长度,而不是字符串。也许你想要LENGTH(CAST(value as string))
  • 除此之外,Kafka 本身规定了默认的最大允许消息大小,但您会在代理/生产者到达 Spark 之前看到该错误

标签: string apache-spark pyspark spark-structured-streaming truncate


【解决方案1】:

类似于控制台截​​断之类的东西。不是 Kafka 或 Spark 问题。

首先我在跑步

# kafka-console-producer.sh --topic dataflow --bootstrap-server localhost:9092

然后将消息粘贴到它的命令行并发生截断。

我跑了

# kafka-console-producer.sh --topic dataflow --bootstrap-server localhost:9092 < row01.json

row01.json 中使用相同的数据,并且它可以在没有截断的情况下工作。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-01-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-09-23
    • 1970-01-01
    相关资源
    最近更新 更多