Scala-Spark消费Kafka Avro消息分Topic生成DataFrame方案
Spark 2.4.4 环境消费Kafka Avro消息分业务实体生成结构化DataFrame实现方案
方案选型
放弃手动初始化KafkaConsumer轮询拉取、手动写文件的实现逻辑,直接使用Spark原生Kafka Batch/Streaming Source对接Kafka,集成Confluent Avro反序列化能力,全程通过DataFrame API完成结构化解析、分实体拆分、存储,不需要手动维护消费者偏移量、文件IO、对象转结构化数据的冗余代码,稳定性和性能适配Spark运行环境。
注意:配置中使用的Schema Registry地址必须和原有原生消费者配置的地址完全一致,否则会触发Avro反序列化报错
第一步:任务依赖配置
提交任务时引入以下版本适配的依赖包(版本匹配Spark 2.4.4 + Scala 2.11 + Confluent 5.3.x,和现有使用的KafkaAvroDeserializer版本保持一致即可):
--packages org.apache.spark:spark-sql-kafka-0-10_2.11:2.4.4,io.confluent:kafka-avro-serializer:5.3.0
第二步:核心实现代码
1. 初始化运行环境与Kafka连接参数
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ import scala.collection.mutable val spark = SparkSession.builder() .appName("KafkaAvro2EntityDF") // 本地调试时打开,集群部署移除该行 //.master("local[*]") .getOrCreate() // 替换为实际环境的配置值 val kafkaParams = Map( "kafka.bootstrap.servers" -> "实际Kafka Broker地址", "subscribe" -> "topic1,topic2,topic3", // 配置需要消费的多Topic列表,逗号分隔 "startingOffsets" -> "earliest", "key.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer", "value.deserializer" -> "io.confluent.kafka.serializers.KafkaAvroDeserializer", "schema.registry.url" -> "实际Schema Registry地址", "specific.avro.reader" -> "false" )
2. 读取Kafka原始数据
// 批量读取数据,如果需要做实时流式消费,把read改成readStream即可 val rawKafkaDF = spark.read .format("kafka") .options(kafkaParams) .load() .select( col("topic").as("source_topic"), col("value").cast("string").as("json_value") )
3. 解析公共外层结构
根据给出的消息样例,所有消息外层固定包含beforeData、headers两个公共字段,业务字段挂在和实体同名的一级Key下:
val commonSchema = StructType(Seq( StructField("beforeData", StringType, nullable = true), StructField("headers", StructType(Seq( StructField("operation", StringType, nullable = true), StructField("changeSequence", StringType, nullable = true), StructField("timestamp", StringType, nullable = true), StructField("externalSchemaId", StringType, nullable = true) )), nullable = true) )) val parsedCommonDF = rawKafkaDF .withColumn("parsed_common", from_json(col("json_value"), commonSchema)) .select( col("source_topic"), col("json_value"), col("parsed_common.headers.operation").as("op_type"), col("parsed_common.headers.externalSchemaId").as("schema_id") )
4. 拆分不同业务实体生成独立DataFrame
// 收集数据中所有出现过的业务实体名称(例如样例中的HML_A_DATA、HML_C_DATA) val entityNames = mutable.Set[String]() rawKafkaDF.select("json_value").as[String].collect().foreach(jsonStr => { val fieldNames = spark.read.json(Seq(jsonStr).toDS).schema.fieldNames fieldNames.foreach(field => { if (!List("beforeData", "headers").contains(field)) entityNames.add(field) }) }) val entityDFMap = mutable.Map[String, org.apache.spark.sql.DataFrame]() entityNames.foreach(entityName => { // 动态提取当前实体的Schema结构 val entityInnerSchema = spark.read .json(rawKafkaDF.filter(col("json_value").contains(s""""$entityName":""")).select("json_value").as[String]) .schema .apply(entityName) .dataType .asInstanceOf[StructType] // 展开业务字段生成二维结构化表 val entityDF = parsedCommonDF .where(col("json_value").contains(s""""$entityName":""")) .withColumn("entity_data", from_json(col("json_value"), StructType(Seq(StructField(entityName, entityInnerSchema, nullable = true))))) .select( col("source_topic"), col("op_type"), col("schema_id"), // 按需追加需要提取的业务字段即可 col(s"entity_data.$entityName.HML_ID").as("hml_id"), col(s"entity_data.$entityName.CREAT_DATETIME").as("creat_datetime"), col(s"entity_data.$entityName.COV_ORDR_TYP_ID").as("cov_order_type_id"), col(s"entity_data.$entityName.COB_STRT_DT").as("cob_start_date"), col(s"entity_data.$entityName.MBR_ID").as("mbr_id"), col(s"entity_data.$entityName.MBR_COV_SEQ_NBR").as("mbr_cov_seq_nbr") ) .dropDuplicates("hml_id", "creat_datetime") entityDFMap.put(entityName, entityDF) })
提示:如果业务实体固定(例如已知只有HML_A_DATA、HML_C_DATA两类),可以跳过动态收集实体名的步骤,直接写死实体列表和对应字段,不需要提前扫描全量数据识别Schema,性能更优。
5. 分实体存储数据
entityDFMap.foreach{ case (entityName, df) => // 控制台打印预览结果 println(s"===== 实体 $entityName 数据预览 =====") df.show(5, false) // 按实体分目录存储,支持本地路径、HDFS、S3等文件系统 df.write .mode("overwrite") // POC测试用overwrite,生产场景根据需求改为append .parquet(s"./output/$entityName/") }
原有FileWriter实现的已知问题
- 逐行写入时反复创建、关闭文件IO流,性能极差,数据量稍大就会触发写入阻塞、文件句柄泄漏
- 直接写入
GenericRecord.toString的输出结果,没有做JSON合法性校验,容易出现格式残缺、转义错误,后续无法直接解析 - 没有内置按业务实体拆分的逻辑,所有实体数据混存到单个文件,后续结构化解析成本极高
- 手动维护原生Kafka消费者实例,无法复用Spark的分布式消费、偏移量持久化、容错重试能力,仅能支撑极小数据量的临时调试,无法适配正式数据处理流程。
内容的提问来源于stack exchange,提问作者Buzz97
相关产品推荐
相关产品推荐

