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
相关产品推荐
相关产品推荐

