You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark Scala如何过滤DataFrame坏记录并将异常记录保存至单独文件

Spark Streaming坏数据分流落地方案

方案1:使用Spark原生坏记录捕获机制

如果是结构化流(Structured Streaming)场景,Spark 2.3+内置支持坏数据重定向能力,解析数据时配置badRecordsPath参数即可自动将格式/类型不匹配的行写入指定路径,同时保留原始数据内容:

// 读取Kafka原始数据
val kafkaDf = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "你的kafka地址:9092")
  .option("subscribe", "你的topic名")
  .load()
  .selectExpr("CAST(value AS STRING) as json_str")

// 解析JSON时开启坏记录自动收集
val parsedDf = spark.read
  .option("badRecordsPath", "/data/spark/bad_records")
  .json(kafkaDf.select("json_str").as[String])

该方案适配大多数解析类错误场景,但无法覆盖写入数据库时触发的约束类错误(如唯一键冲突、字段长度超限)。

方案2:自定义校验+分流写入(全场景兼容)

该方案可控性最高,同时支持捕获解析、校验、入库全链路的错误数据,还会自动保留Kafka元信息,无需后续排查偏移量:

  • 第一步:给原始数据追加is_valid、error_msg、kafka_topic、kafka_partition、kafka_offset辅助字段
  • 第二步:用filter算子将DataFrame拆分为正常流、错误流两个独立分支
  • 第三步:正常流走原有的数据库写入逻辑,错误流将原始内容、错误信息、Kafka元信息统一写入指定存储路径

代码示例:

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._

// 定义目标表的Schema
val targetSchema = StructType(Seq(
  StructField("id", IntegerType, nullable = false),
  StructField("user_name", StringType, nullable = false),
  StructField("create_time", TimestampType, nullable = false)
))

// 读取Kafka数据并保留元信息
val rawDf = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "你的kafka地址:9092")
  .option("subscribe", "你的topic名")
  .load()
  .select(
    col("value").cast(StringType).alias("raw_data"),
    col("topic").alias("kafka_topic"),
    col("partition").alias("kafka_partition"),
    col("offset").alias("kafka_offset")
  )

// 逐行校验标记错误数据
val validatedDf = rawDf
  .withColumn("parsed_data", from_json(col("raw_data"), targetSchema))
  .withColumn("is_valid", 
    col("parsed_data").isNotNull 
    && col("parsed_data.id").isNotNull 
    && col("parsed_data.create_time").isNotNull
  )
  .withColumn("error_msg", when(col("is_valid"), lit(null)).otherwise(lit("数据格式/类型不匹配")))

// 拆分数据流
val validDf = validatedDf.filter(col("is_valid")).select("parsed_data.*")
val invalidDf = validatedDf.filter(!col("is_valid")).select("raw_data", "kafka_topic", "kafka_partition", "kafka_offset", "error_msg")

// 正常数据写入数据库
val validWriteTask = validDf.writeStream
  .format("jdbc")
  .option("url", "jdbc:mysql://你的数据库地址/库名")
  .option("dbtable", "目标表名")
  .option("user", "数据库账号")
  .option("password", "数据库密码")
  .option("checkpointLocation", "/data/spckpt/valid_stream")
  .start()

// 错误数据写入文件
val invalidWriteTask = invalidDf.writeStream
  .format("json")
  .option("path", "/data/spark/error_records")
  .option("checkpointLocation", "/data/spckpt/invalid_stream")
  .start()

spark.streams.awaitAnyTermination()

注意事项

  • 两个流的checkpointLocation必须分开配置,避免状态冲突导致作业异常
  • 错误存储路径建议按天分区,方便后续按时间维度回溯问题
  • 如果需要捕获入库阶段的约束错误,可将写入逻辑替换为foreachBatch算子,在批次内对JDBC操作做try-catch捕获异常行单独落盘

内容的提问来源于stack exchange,提问作者mt_leo

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.05 12:12:03