spark-streaming-kafka-0-10与spark-sql-kafka-0-10的区别及批处理选型
Spark Kafka 相关库区别及批处理场景选型
我的批处理场景代码
我需要读取Parquet文件并写入Kafka,对应的Spark代码如下:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.struct import org.apache.spark.sql.functions.to_json object IngestFromS3ToKafka { def main(args: Array[String]): Unit = { val spark: SparkSession = SparkSession .builder() .master("local[*]") .appName("ingest-from-s3-to-kafka") .config("spark.ui.port", "4040") .getOrCreate() val filePath = "s3a://my-bucket/my.parquet" spark.read.parquet(filePath) .select(to_json(struct("*")).alias("value")) .write .format("kafka") .option("kafka.bootstrap.servers", "hm-kafka-kafka-bootstrap.hm-kafka.svc:9092") .option("topic", "my-topic") .save() spark.stop() } }
疑问
根据《Structured Streaming + Kafka Integration Guide》文档,我了解到spark-sql-kafka-0-10库支持批处理和流处理,但发现两个相关库:
spark-streaming-kafka-0-10:Spark Integration For Kafka 0.10spark-sql-kafka-0-10:Kafka 0.10+ Source For Structured Streaming
我的场景是批处理而非流处理,但这两个库的名称和描述都与流处理相关,想知道这两个库的区别是什么?是否有相关文档说明?
解答
两个库的核心区别
spark-streaming-kafka-0-10- 属于Spark早期的**DStream API(Spark Streaming)**体系,是专门为传统Spark Streaming流处理设计的组件,基于RDD模型开发。
- 仅支持Spark Streaming的流处理场景,完全不兼容你代码里用到的SparkSession/DataFrame/SQL API,无法实现批处理写入Kafka的操作。
spark-sql-kafka-0-10- 属于**Structured Streaming(Spark SQL)**生态,是Spark官方推荐的与Kafka交互的组件。
- 虽然名称和描述关联Structured Streaming,但它同时支持流处理(
readStream/writeStream)和批处理(read/write)两种场景,你代码中用write.format("kafka")写入Kafka的逻辑,正是依赖这个库实现的。
文档依据
Spark官方文档明确说明:spark-sql-kafka-0-10是唯一支持通过Structured SQL/DataFrame API与Kafka进行交互的库,覆盖批处理和流处理场景;而spark-streaming-kafka-0-10仅服务于旧版Spark Streaming(DStream)的流处理需求,与当前主流的SparkSession API体系不兼容。
内容的提问来源于stack exchange,提问作者Hongbo Miao
相关产品推荐
相关产品推荐

