[ 
https://issues.apache.org/jira/browse/FLINK-40232?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Natea Eshetu Beshada reassigned FLINK-40232:
--------------------------------------------

    Assignee: Iman Rastkhadiv

> NullPointerException in OutputConversionOperator.processElement when rowtime 
> attribute is null
> ----------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40232
>                 URL: https://issues.apache.org/jira/browse/FLINK-40232
>             Project: Flink
>          Issue Type: Bug
>          Components: Table SQL / API, Table SQL / Planner
>    Affects Versions: 1.20.5, 2.1.3
>         Environment: Flink 1.15+ (verified on release-1.15, release-2.1, 
> master)
>            Reporter: Fan Deng
>            Assignee: Iman Rastkhadiv
>            Priority: Major
>
> h2. Problem
>   {{OutputConversionOperator.processElement}} (line 74) throws 
> {{NullPointerException}} when a rowtime attribute is {{null}} and 
> {{consumeRowtimeMetadata=true}}. The operator calls
>   {{rowData.getTimestamp(...).getMillisecond()}} without null check.
>   h2. Root Cause
>   File: 
> {{flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/OutputConversionOperator.java}}
>   {code:java}
>   @Override
>   public void processElement(StreamRecord<RowData> element) throws Exception {
>       final RowData rowData = element.getValue();
>   if (consumeRowtimeMetadata) {
>       // timestamp is TIMESTAMP_LTZ
>       final long rowtime = rowData.getTimestamp(rowData.getArity() - 1, 
> 3).getMillisecond();  // ← line 74 NPE
>       outRecord.setTimestamp(rowtime);
>   } else if (rowtimeIndex != -1) {
>       // timestamp might be TIMESTAMP or TIMESTAMP_LTZ
>       final long rowtime = rowData.getTimestamp(rowtimeIndex, 
> 3).getMillisecond();  // ← line 78 same issue
>       outRecord.setTimestamp(rowtime);
>   }
>   ...
>   }
>   {code}
>   The code calls {{rowData.getTimestamp(...)}} and immediately invokes 
> {{.getMillisecond()}} on the result without null check. When the rowtime 
> field is {{null}}, {{getTimestamp()}} returns Java {{null}} for certain
>   {{RowData}} implementations, causing {{NullPointerException}}.
>   h2. Why NPE Occurs with BoxedWrapperRowData
>   Calc codegen produces {{BoxedWrapperRowData}} as output (see 
> {{CalcCodeGenerator.scala:61}}):
>   {code:scala}
>   // CalcCodeGenerator.scala:61
>   classOf[BoxedWrapperRowData],
>   {code}
>   {{BoxedWrapperRowData.getTimestamp}} returns Java {{null}} for null fields 
> (unlike {{BinaryRowData}} which returns {{TimestampData.fromEpochMillis(0)}}):
>   {code:java}
>   // BoxedWrapperRowData.java:115-116
>   public TimestampData getTimestamp(int pos, int precision) {
>       return (TimestampData) this.fields[pos];  // fields[pos]==null → 
> returns Java null
>   }
>   {code}
>   ||RowData implementation||{{getTimestamp(pos, precision)}} for null 
> field||{{.getMillisecond()}}||
>   |{{BinaryRowData}}|{{TimestampData.fromEpochMillis(0)}} (non-null)|Returns 
> 0, no NPE|
>   |{{BoxedWrapperRowData}}|Java {{null}}|NPE|
>   |{{GenericRowData}}|Java {{null}}|NPE|
>   h2. Trigger Conditions
>   The NPE requires all of the following:
>   {{consumeRowtimeMetadata=true}}: The sink schema declares a rowtime 
> metadata column (e.g., via {{Schema.columnByMetadata("rowtime", ...)}})
>   Null rowtime field: A record with null rowtime reaches 
> {{OutputConversionOperator}}
>   {{BoxedWrapperRowData}} passed directly: Calc codegen output 
> ({{BoxedWrapperRowData}}) is passed directly to {{OutputConversionOperator}} 
> without serialization (operator chaining enabled, object reuse enabled)
>   No upstream {{WatermarkAssignerOperator}}: Source performs watermark 
> pushdown (e.g., Kafka source implements {{SupportsWatermarkPushDown}}), so 
> the planner does not insert an independent
>   {{WatermarkAssignerOperator}} that would otherwise reject null rowtime 
> earlier
>   h2. Reproduction
>   h3. End-to-end MiniCluster test
>   Verified on release-1.15:
>   {code:java}
>   // Source with watermark pushdown (no independent WatermarkAssignerOperator)
>   tenv.executeSql(
>       "CREATE TABLE src (\n"
>           + "  f0 STRING,\n"
>           + "  rtime TIMESTAMP(3),\n"
>           + "  WATERMARK FOR rtime AS rtime - INTERVAL '10' SECOND\n"
>           + ") WITH (\n"
>           + "  'connector' = 'pushdown-watermark-values',\n"
>           + "  'data-id' = '...'  // emits Row with null rtime\n"
>           + ")");
>   Table table = tenv.from("src");
>   // Sink schema with rowtime metadata column → consumeRowtimeMetadata=true
>   Schema sinkSchema = Schema.newBuilder()
>       .column("f0", DataTypes.STRING())
>       .columnByMetadata("rtime", DataTypes.TIMESTAMP_LTZ(3).notNull(), 
> "rowtime")
>       .build();
>   // toChangelogStream triggers 
> OutputConversionOperator(consumeRowtimeMetadata=true)
>   // Calc codegen produces BoxedWrapperRowData → getTimestamp returns null → 
> NPE
>   tenv.toChangelogStream(table, sinkSchema, 
> ChangelogMode.insertOnly()).print();
>   env.execute("repro");  // → NPE at 
> OutputConversionOperator.processElement:74
>   {code}
>   h3. Stack trace
>   {code}
>   java.lang.NullPointerException
>       at 
> org.apache.flink.table.runtime.operators.sink.OutputConversionOperator.processElement(OutputConversionOperator.java:74)
>       at 
> org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.pushToOperator(CopyingChainingOutput.java:82)
>       ...
>   Caused by: java.lang.NullPointerException
>       at 
> org.apache.flink.table.runtime.operators.sink.OutputConversionOperator.processElement(OutputConversionOperator.java:74)
>   {code}
>   h2. Why disable-operator-chaining works as workaround
>   When operator chaining is disabled, the chain boundary forces 
> {{RowDataSerializer}} serialization/deserialization. The deserialized row 
> becomes {{BinaryRowData}}, whose {{getTimestamp}} is null-tolerant (returns
>   {{TimestampData.fromEpochMillis(0)}} for null fields). This is not a fix — 
> it's a workaround that masks the underlying bug.
>   h2. Impact
>   - Affects: release-1.15, release-2.1, master (verified byte-identical via 
> {{git diff}})
>   - Introduced in: FLINK-19980 (Timo Walther, 2021-04-13)
>   - Triggered in production: Oceanus (Tencent's Flink platform) — Kafka 
> source with source-watermark pushdown + {{toChangelogStream}} with rowtime 
> metadata column
>   - No existing JIRA: This bug has not been reported to the community for 5+ 
> years
>   h2. Suggested Fix
>   Add null check in {{OutputConversionOperator.processElement}}:
>   {code:java}
>   @Override
>   public void processElement(StreamRecord<RowData> element) throws Exception {
>       final RowData rowData = element.getValue();
>   if (consumeRowtimeMetadata) {
>       // timestamp is TIMESTAMP_LTZ
>       TimestampData ts = rowData.getTimestamp(rowData.getArity() - 1, 3);
>       if (ts != null) {
>           final long rowtime = ts.getMillisecond();
>           outRecord.setTimestamp(rowtime);
>       }
>   } else if (rowtimeIndex != -1) {
>       // timestamp might be TIMESTAMP or TIMESTAMP_LTZ
>       TimestampData ts = rowData.getTimestamp(rowtimeIndex, 3);
>       if (ts != null) {
>           final long rowtime = ts.getMillisecond();
>           outRecord.setTimestamp(rowtime);
>       }
>   }
>   ...
>   }
>   {code}



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to