Scala读取含无引号时间戳的JSON/Ion日志报错,求解决方案
问题场景
现有如下非标准JSON日志,其中time_of_initial_request字段的时间戳值未加双引号,字段名也无引号:
{post_handler_type:"trace", default_marketplace:"prod_iad", request_id:"5MTB3656X0JRL38R6LJY", rest_uri:"productApi/feature", client_logging_id:"PredictionMetadata", time_of_initial_request:2023-05-24T15:00:00.577Z }
使用Spark的allowUnquotedFieldNames参数读取时,因无引号的时间戳被识别为非法JSON语法,导致损坏记录报错。测试发现给该时间戳添加双引号后可正常解析,但无法修改原日志格式,且业务上不需要time_of_initial_request字段,需在Scala环境下解决该问题。
原测试代码:
val JsonRaw = spark.read.option("allowUnquotedFieldNames","true") .option("inferSchema", "true") .json("s3://my-bucket/prediction-test/2023-05-24/") val df = JsonRaw.registerTempTable("prediction") val df2 = spark.sql("select request_id from prediction") df2.show()
解决方案
方案1:文本预处理移除异常字段
先将文件按文本行读取,通过正则表达式移除time_of_initial_request字段及其值,再将清理后的文本转为JSON格式解析:
import org.apache.spark.sql.functions.regexp_replace // 按文本读取原始日志 val rawLines = spark.read.text("s3://my-bucket/prediction-test/2023-05-24/") // 正则替换掉time_of_initial_request字段(处理字段在中间或末尾的情况) val cleanedLines = rawLines .withColumn("value", regexp_replace($"value", """time_of_initial_request:\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}\.\d{3}Z,?""", "")) // 处理字段在末尾时遗留的逗号 .withColumn("value", regexp_replace($"value", ",\\s*}", "}")) // 读取清理后的合法JSON val jsonDF = spark.read .option("allowUnquotedFieldNames", "true") .json(cleanedLines) // 提取需要的字段 val resultDF = jsonDF.select("request_id") resultDF.show()
优点:实现简单,性能较高,适合日志格式固定的场景。
方案2:使用Jackson自定义解析规则
利用Jackson的灵活配置,支持无引号字段名并忽略解析异常,通过UDF逐行解析日志并提取目标字段:
import com.fasterxml.jackson.databind.json.JsonMapper import com.fasterxml.jackson.module.scala.DefaultScalaModule import org.apache.spark.sql.functions.udf // 配置Jackson解析器:支持无引号字段名、忽略未知属性、容错解析 val mapper = JsonMapper.builder() .enable(com.fasterxml.jackson.databind.DeserializationFeature.ALLOW_UNQUOTED_FIELD_NAMES) .disable(com.fasterxml.jackson.databind.DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES) .disable(com.fasterxml.jackson.databind.DeserializationFeature.FAIL_ON_INVALID_SUBTYPE) .addModule(DefaultScalaModule) .build() // 定义UDF,仅提取需要的request_id字段,避免解析整个对象 val extractRequestIdUdf = udf((line: String) => { try { val jsonNode = mapper.readTree(line) Option(jsonNode.get("request_id")).map(_.asText()) } catch { case _: Exception => None } }) // 读取文本并解析提取目标字段 val resultDF = spark.read.text("s3://my-bucket/prediction-test/2023-05-24/") .select(extractRequestIdUdf($"value").alias("request_id")) .filter($"request_id".isNotNull) resultDF.show()
优点:容错性更强,能处理更多格式异常的日志,适合复杂多变的日志场景。
内容的提问来源于stack exchange,提问作者Ankur
相关产品推荐
相关产品推荐

