[
https://issues.apache.org/jira/browse/FLINK-10119?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Flink Jira Bot updated FLINK-10119:
-----------------------------------
Labels: auto-deprioritized-major auto-unassigned pull-request-available
stale-minor (was: auto-deprioritized-major auto-unassigned
pull-request-available)
I am the [Flink Jira Bot|https://github.com/apache/flink-jira-bot/] and I help
the community manage its development. I see this issues has been marked as
Minor but is unassigned and neither itself nor its Sub-Tasks have been updated
for 180 days. I have gone ahead and marked it "stale-minor". If this ticket is
still Minor, please either assign yourself or give an update. Afterwards,
please remove the label or in 7 days the issue will be deprioritized.
> JsonRowDeserializationSchema deserialize kafka message
> ------------------------------------------------------
>
> Key: FLINK-10119
> URL: https://issues.apache.org/jira/browse/FLINK-10119
> Project: Flink
> Issue Type: Bug
> Components: Connectors / Kafka
> Affects Versions: 1.5.1
> Environment: 无
> Reporter: sean.miao
> Priority: Minor
> Labels: auto-deprioritized-major, auto-unassigned,
> pull-request-available, stale-minor
>
> Recently, we are using Kafka010JsonTableSource to process kafka's json
> messages.We turned on checkpoint and auto-restart strategy .
> We found that as long as the format of a message is not json, it will cause
> the job to not be pulled up. Of course, this is to ensure that only once
> processing or at least once processing, but the resulting application is not
> available and has a greater impact on us.
> the code is :
> class : JsonRowDeserializationSchema
> function :
> @Override
> public Row deserialize(byte[] message) throws IOException {
> try
> { final JsonNode root = objectMapper.readTree(message); return
> convertRow(root, (RowTypeInfo) typeInfo); }
> catch (Throwable t)
> { throw new IOException("Failed to deserialize JSON object.", t); }
> }
> now ,i change it to :
> public Row deserialize(byte[] message) throws IOException {
> try
> { JsonNode root = this.objectMapper.readTree(message); return
> this.convertRow(root, (RowTypeInfo)this.typeInfo); }
> catch (Throwable var4) {
> message = this.objectMapper.writeValueAsBytes("{}");
> JsonNode root = this.objectMapper.readTree(message);
> return this.convertRow(root, (RowTypeInfo)this.typeInfo);
> }
> }
>
> I think that data format errors are inevitable during network transmission,
> so can we add a new column to the table for the wrong data format? like spark
> sql does。
>
--
This message was sent by Atlassian Jira
(v8.20.1#820001)