要从 Kafka 读取数据并将其以 Parquet 格式写入 HDFS,使用 Spark Batch 作业而不是流式传输,您可以使用 Spark Structured Streaming。
Structured Streaming 是基于 Spark SQL 引擎构建的可扩展和容错流处理引擎。您可以像表达对静态数据的批处理计算一样表达您的流计算。 Spark SQL 引擎将负责以增量和连续的方式运行它,并随着流数据的不断到达而更新最终结果。您可以使用 Scala、Java、Python 或 R 中的 Dataset/DataFrame API 来表示流聚合、事件时间窗口、流到批处理连接等。计算在同一个优化的 Spark SQL 引擎上执行。最后,系统通过检查点和预写日志确保端到端的精确一次容错保证。简而言之,结构化流式处理提供了快速、可扩展、容错、端到端的一次性流处理,而用户无需对流式处理进行推理。
它带有 Kafka 作为内置 Source,即我们可以从 Kafka 轮询数据。它与 Kafka 代理版本 0.10.0 或更高版本兼容。
为了以批处理模式从 Kafka 中提取数据,您可以为定义的偏移范围创建 Dataset/DataFrame。
// Subscribe to 1 topic defaults to the earliest and latest offsets
val df = spark
.read
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,host2:port2")
.option("subscribe", "topic1")
.load()
df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
.as[(String, String)]
// Subscribe to multiple topics, specifying explicit Kafka offsets
val df = spark
.read
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,host2:port2")
.option("subscribe", "topic1,topic2")
.option("startingOffsets", """{"topic1":{"0":23,"1":-2},"topic2":{"0":-2}}""")
.option("endingOffsets", """{"topic1":{"0":50,"1":-1},"topic2":{"0":-1}}""")
.load()
df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
.as[(String, String)]
// Subscribe to a pattern, at the earliest and latest offsets
val df = spark
.read
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,host2:port2")
.option("subscribePattern", "topic.*")
.option("startingOffsets", "earliest")
.option("endingOffsets", "latest")
.load()
df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
.as[(String, String)]
源中的每一行都有以下架构:
| Column | Type |
|:-----------------|--------------:|
| key | binary |
| value | binary |
| topic | string |
| partition | int |
| offset | long |
| timestamp | long |
| timestampType | int |
现在,要将数据以 parquet 格式写入 HDFS,可以编写以下代码:
df.write.parquet("hdfs://data.parquet")
有关 Spark Structured Streaming + Kafka 的更多信息,请参阅以下指南 - Kafka Integration Guide
希望对你有帮助!