【发布时间】: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