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

Flink Table API消费Kafka JSON消息反序列化失败求助

我用Flink Table API开发了一个应用,从自建Kafka Topic读取Kaggle的YouTube统计数据。Confluent UI显示Topic已收到生产者消息且格式正常,单条消息JSON也通过在线工具验证有效。但应用在将数据转换后写入另一个Kafka Topic时,出现源JSON反序列化失败;开启value.json.ignore-parse-errors = 'true'后异常消失,但没有数据输出到目标Topic。

异常信息

Caused by: java.io.IOException: Failed to deserialize consumer record ConsumerRecord(topic = youtube_stats, partition = 0, leaderEpoch = 0, offset = 0, CreateTime = 1675701408811, serialized key size = 12, serialized value size = 1955, headers = RecordHeaders(headers = [], isReadOnly = false), key = [B@3dd62069, value = [B@2aa4839f).
        at org.apache.flink.connector.kafka.source.reader.deserializer.KafkaDeserializationSchemaWrapper.deserialize(KafkaDeserializationSchemaWrapper.java:57)
        at org.apache.flink.connector.kafka.source.reader.KafkaRecordEmitter.emitRecord(KafkaRecordEmitter.java:53)
        ... 14 more
Caused by: java.io.IOException: Failed to deserialize JSON '{"video_id": "2kyS6SvSYSE", "trending_date": "17.14.11", "title": "WE WANT TO TALK ABOUT OUR MARRIAGE", "channel_title": "CaseyNeistat", "category_id": 22, "publish_time": "2017-11-13T17:13:01.000Z", "tags": "SHANtell martin", "views": 748374, "likes": 57527, "dislikes": 2966, "comment_count": 15954, "thumbnail_link": "https://i.ytimg.com/vi/2kyS6SvSYSE/default.jpg", "comments_disabled": false, "ratings_disabled": false, "video_error_or_removed": false, "description": "SHANTELL'S CHANNEL - https://www.youtube.com/shantellmartin\nCANDICE - https://www.lovebilly.com\n\nfilmed this video in 4k on this -- http://amzn.to/2sTDnRZ\nwith this lens -- http://amzn.to/2rUJOmD\nbig drone - http://someurl.com/h4ft3oy\nOTHER GEAR ---  http://amzn.to/2o3GLX5\nSony CAMERA http://amzn.to/2nOBmnv\nOLD CAMERA; http://amzn.to/2o2cQBT\nMAIN LENS; http://amzn.to/2od5gBJ\nBIG SONY CAMERA; http://amzn.to/2nrdJRO\nBIG Canon CAMERA; http://someurl.com/jn4q4vz\nBENDY TRIPOD THING; http://someurl.com/gw3ylz2\nYOU NEED THIS FOR THE BENDY TRIPOD; http://someurl.com/j8mzzua\nWIDE LENS; http://someurl.com/jkfcm8t\nMORE EXPENSIVE WIDE LENS; http://someurl.com/zrdgtou\nSMALL CAMERA; http://someurl.com/hrrzhor\nMICROPHONE; http://someurl.com/zefm4jy\nOTHER MICROPHONE; http://someurl.com/jxgpj86\nOLD DRONE (cheaper but still great);http://someurl.com/zcfmnmd\n\nfollow me; on http://instagram.com/caseyneistat\non https://www.facebook.com/cneistat\non https://twitter.com/CaseyNeistat\n\namazing intro song by https://soundcloud.com/discoteeth\n\nad disclosure.  THIS IS NOT AN AD.  not selling or promoting anything.  but samsung did produce the Shantell Video as a 'GALAXY PROJECT' which is an initiative that enables creators like Shantell and me to make projects we might otherwise not have the opportunity to make.  hope that's clear.  if not ask in the comments and i'll answer any specifics.", "event_time": "2023-02-06 16:36:48"}'.
        at org.apache.flink.formats.json.JsonRowDataDeserializationSchema.deserialize(JsonRowDataDeserializationSchema.java:122)
        at org.apache.flink.formats.json.JsonRowDataDeserializationSchema.deserialize(JsonRowDataDeserializationSchema.java:51)
        at org.apache.flink.api.common.serialization.DeserializationSchema.deserialize(DeserializationSchema.java:82)
        at org.apache.flink.streaming.connectors.kafka.table.DynamicKafkaDeserializationSchema.deserialize(DynamicKafkaDeserializationSchema.java:113)
        at org.apache.flink.connector.kafka.source.reader.deserializer.KafkaDeserializationSchemaWrapper.deserialize(KafkaDeserializationSchemaWrapper.java:54)
        ... 15 more
Caused by: java.io.CharConversionException: Invalid UTF-32 character 0x17a2276 (above 0x0010ffff) at char #1, byte #7)
        at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.io.UTF32Reader.reportInvalid(UTF32Reader.java:195)
        at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.io.UTF32Reader.read(UTF32Reader.java:158)
        at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.json.ReaderBasedJsonParser._loadMore(ReaderBasedJsonParser.java:255)
        at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.json.ReaderBasedJsonParser._skipWSOrEnd(ReaderBasedJsonParser.java:2389)
        at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.json.ReaderBasedJsonParser.nextToken(ReaderBasedJsonParser.java:677)
        at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper._readTreeAndClose(ObjectMapper.java:4622)
        at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper.readTree(ObjectMapper.java:3056)
        at org.apache.flink.formats.json.JsonRowDataDeserializationSchema.deserializeToJsonNode(JsonRowDataDeserializationSchema.java:127)
        at org.apache.flink.formats.json.JsonRowDataDeserializationSchema.deserialize(JsonRowDataDeserializationSchema.java:116)
        ... 19 more

相关代码片段

fields = [
            ("video_id", "STRING"),
            ("trending_date", "STRING"),
            ("title", "STRING"),
            ("channel_title", "STRING"),
            ("category_id", "STRING"),
            ("publish_time", "STRING"),
            ("tags", "STRING"),
            ("`views`", "BIGINT"),
            ("likes", "BIGINT"),
            ("dislikes", "BIGINT"),
            ("comment_count", "BIGINT"),
            ("thumbnail_link", "STRING"),
            ("comments_disabled", "BOOLEAN"),
            ("ratings_disabled", "BOOLEAN"),
            ("video_error_or_removed", "BOOLEAN"),
            ("description", "STRING"),
            ("event_time", "TIMESTAMP(3)"),
            ("WATERMARK FOR event_time AS event_time - INTERVAL '10' second", "")
        ]

table_ddl = """
CREATE TABLE {kwargs["table_name"]} (
                {','.join([f"{ft[0]} {ft[1]}" for ft in fields])}
            ) WITH (
                'connector' = 'kafka',
                'topic' = '{kwargs["topic"]}',
                'properties.bootstrap.servers' = '{kwargs["brokers"]}',
                'scan.startup.mode' = 'earliest-offset',
                'value.format' = '{kwargs["value_format"]}',
                'value.json.fail-on-missing-field' = 'false',
                'value.json.ignore-parse-errors' = 'true'
            )
"""
table_env.execute_sql(table_ddl)

sink_table = # pretty similar to the other

tumble_window = (
            Tumble.over(lit(self.agg_every_n_secs).seconds)
            .on(col("event_time"))
            .alias("w")
        )
stats_table = (
            table_env.from_path(table_name).window(tumble_window)
            .group_by(col("w"), col("channel_title"))
            .select(
                col("channel_title"),
                col("likes").max.alias("likes"),
                col("dislikes").max.alias("dislikes"),
                (
                    col("likes").cast(DataTypes.FLOAT()).max
                    / col("dislikes").if_null(1).cast(DataTypes.FLOAT()).max
                ).alias("like_dislike_ratio"),
                col("views").sum.alias("total_views"),
                col("w").end.cast(DataTypes.TIMESTAMP(3)).alias("proctime"),
            )
        )
stats_table.execute_insert(sink_table_name).wait()

解决方案

  • 修复UTF编码问题:错误栈核心是Invalid UTF-32 character 0x17a2276,说明生产者发送的消息编码不符合UTF-8标准。检查生产者端的消息生成逻辑,确保输出的是UTF-8格式的JSON字符串;如果生产者使用了其他编码,在Flink表DDL中添加'value.json.charset' = '对应编码'指定正确字符集。
  • 排查字段类型不匹配:临时关闭value.json.ignore-parse-errors,只保留value.json.fail-on-missing-field = 'false',重新运行后查看具体的类型转换错误。比如JSON中category_id是数字类型,但表定义为STRING,虽然Flink支持自动转换,但其他字段的类型不兼容可能导致数据被过滤。
  • 验证窗口聚合逻辑:先去掉窗口聚合,直接将源表数据写入目标Topic,确认是否有数据输出。如果有数据,再逐步添加聚合逻辑,检查agg_every_n_secs是否设置过大导致窗口未触发,或者event_time格式错误导致水印不推进、窗口无法关闭。
  • 查看Flink日志:开启ignore-parse-errors后,Flink会跳过错误数据并在TaskManager日志中记录跳过的数量,通过日志确认是所有数据都被跳过,还是部分数据解析成功但聚合无输出。
  • 重置Kafka消费者offset:如果Topic中的消息已被消费过,在DDL中添加'properties.group.id' = '新消费者组ID',或者重置原有消费者组的offset,确保Flink能重新读取Topic中的数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 20:30:56