这是一种可能的方法:
在开始流式传输之前,从 Kafka 获取一小批数据
从小批量推断架构
使用提取的架构开始流式传输数据。
下面的伪代码说明了这种方法。
第 1 步:
从 Kafka 中提取一小批(两条记录),
val smallBatch = spark.read.format("kafka")
.option("kafka.bootstrap.servers", "node:9092")
.option("subscribe", "topicName")
.option("startingOffsets", "earliest")
.option("endingOffsets", """{"topicName":{"0":2}}""")
.load()
.selectExpr("CAST(value AS STRING) as STRING").as[String].toDF()
第 2 步:
将小批量写入文件:
smallBatch.write.mode("overwrite").format("text").save("/batch")
此命令将小批量写入 hdfs 目录 /batch。它创建的文件的名称是 part-xyz*。因此,您首先需要使用 hadoop FileSystem 命令重命名文件(参见 org.apache.hadoop.fs._ 和 org.apache.hadoop.conf.Configuration,这是一个示例 https://stackoverflow.com/a/41990859),然后将文件读取为 json:
val smallBatchSchema = spark.read.json("/batch/batchName.txt").schema
这里,batchName.txt 是文件的新名称,smallBatchSchema 包含从小批量推断的架构。
最后,您可以按如下方式流式传输数据(第 3 步):
val inputDf = spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", "node:9092")
.option("subscribe", "topicName")
.option("startingOffsets", "earliest")
.load()
val dataDf = inputDf.selectExpr("CAST(value AS STRING) as json")
.select( from_json($"json", schema=smallBatchSchema).as("data"))
.select("data.*")
希望这会有所帮助!