[
https://issues.apache.org/jira/browse/FLINK-40479?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Rango closed FLINK-40479.
-------------------------
Resolution: Not A Bug
2.3 never promised this parameter. The 2.3 docs do not list error_handling, and
the 2.3 release notes call invalid_op_handling future work (FLINK-39430, Fix
Version 2.4.0).
On 2.3.0 the documented FROM_CHANGELOG path works. Unmapped op codes fail at
runtime, which is the intended default after FLINK-39495. Passing
error_handling is ignored and falls back to that default; it does not corrupt
results.
Closing as not a 2.3 defect. Users who need SKIP can wait for 2.4.
> FROM_CHANGELOG accepts error_handling but still fails the job on unmapped op
> codes
> ----------------------------------------------------------------------------------
>
> Key: FLINK-40479
> URL: https://issues.apache.org/jira/browse/FLINK-40479
> Project: Flink
> Issue Type: Bug
> Components: Table SQL / API
> Affects Versions: 2.3.0
> Environment: Apache Flink 2.3.0 (docker image
> flink:2.3.0-scala_2.12-java17), SQL Client, streaming mode
> Reporter: Rango
> Priority: Major
>
> Flink 2.3.0 ({{flink:2.3.0-scala_2.12-java17}}).
> The 2.3 docs for {{FROM_CHANGELOG}} do not list an {{error_handling}}
> parameter, and the 2.3 release notes say {{invalid_op_handling}} is future
> work. The current master docs do list {{error_handling => 'FAIL' | 'SKIP'}}.
> On 2.3.0 the planner accepts {{error_handling}} anyway. The job is submitted,
> but an unmapped op code still throws and the job restarts:
> {code}
> org.apache.flink.table.api.TableRuntimeException:
> Received invalid op code '6'. Defined op codes are:
> [UPDATE_AFTER, UPDATE_BEFORE, INSERT, DELETE]
> at
> org.apache.flink.table.runtime.functions.ptf.FromChangelogFunction.eval(FromChangelogFunction.java:146)
> {code}
> The same happens with {{error_handling => 'IGNORE'}} and with the
> undocumented alias {{op_col}}.
> Expected:
> * either reject unknown / unimplemented arguments at planning time
> * or honor SKIP and drop the bad row instead of failing the task
> h3. Repro (streaming)
> {code:sql}
> CREATE TABLE cdc_raw (
> op STRING,
> id INT,
> name STRING
> ) WITH (
> 'connector' = 'datagen',
> 'number-of-rows' = '8',
> 'fields.op.kind' = 'random',
> 'fields.op.length' = '1'
> );
> CREATE TABLE out_print (id INT, name STRING) WITH ('connector' = 'print');
> INSERT INTO out_print
> SELECT * FROM TABLE(
> FROM_CHANGELOG(
> input => TABLE cdc_raw,
> op => DESCRIPTOR(op),
> error_handling => 'SKIP'
> )
> );
> {code}
> h3. Control
> The same pipeline with VALUES
> {code:sql}
> ('INSERT', 1, 'a'),
> ('UPDATE_BEFORE', 1, 'a'),
> ('UPDATE_AFTER', 1, 'b'),
> ('DELETE', 1, 'b')
> {code}
> finishes successfully. So {{FROM_CHANGELOG}} itself works; only the
> error-handling path is broken / ignored.
> Related: FLINK-39430 (error_handling added), FLINK-39258 (FLIP-564).
--
This message was sent by Atlassian Jira
(v8.20.10#820010)