[
https://issues.apache.org/jira/browse/FLINK-20722?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=17253485#comment-17253485
]
Rui Li edited comment on FLINK-20722 at 12/22/20, 1:23 PM:
-----------------------------------------------------------
Execution plan is:
{noformat}
== Abstract Syntax Tree ==
LogicalSink(table=[test-catalog.default.dest], fields=[key, val])
+- LogicalProject(key=[$0], val=[$1])
+- LogicalUnion(all=[true])
:- LogicalProject(key=[$0], val=[$3])
: +- LogicalJoin(condition=[=($0, $2)], joinType=[left])
: :- LogicalTableScan(table=[[test-catalog, default, src2]])
: +- LogicalTableScan(table=[[test-catalog, default, src1]])
+- LogicalProject(key=[$0], val=[$3])
+- LogicalJoin(condition=[=($0, $2)], joinType=[left])
:- LogicalTableScan(table=[[test-catalog, default, src2]])
+- LogicalTableScan(table=[[test-catalog, default, src1]])
== Optimized Physical Plan ==
Sink(table=[test-catalog.default.dest], fields=[key, val])
+- Union(all=[true], union=[key, val])
:- Calc(select=[key, val])
: +- HashJoin(joinType=[LeftOuterJoin], where=[=(key, key0)], select=[key,
key0, val], build=[left])
: :- Exchange(distribution=[hash[key]])
: : +- TableSourceScan(table=[[test-catalog, default, src2,
project=[key]]], fields=[key])
: +- Exchange(distribution=[hash[key]])
: +- TableSourceScan(table=[[test-catalog, default, src1]],
fields=[key, val])
+- Calc(select=[key, val])
+- HashJoin(joinType=[LeftOuterJoin], where=[=(key, key0)], select=[key,
key0, val], build=[left])
:- Exchange(distribution=[hash[key]])
: +- TableSourceScan(table=[[test-catalog, default, src2,
project=[key]]], fields=[key])
+- Exchange(distribution=[hash[key]])
+- TableSourceScan(table=[[test-catalog, default, src1]],
fields=[key, val])
== Optimized Execution Plan ==
Sink(table=[test-catalog.default.dest], fields=[key, val])
+- MultipleInput(readOrder=[0,1], members=[\nUnion(all=[true], union=[key,
val])\n:- Calc(select=[key, val])(reuse_id=[1])\n: +-
HashJoin(joinType=[LeftOuterJoin], where=[(key = key0)], select=[key, key0,
val], build=[left])\n: :- [#1] Exchange(distribution=[hash[key]])\n: +-
[#2] Exchange(distribution=[hash[key]])\n+- Reused(reference_id=[1])\n])
:- Exchange(distribution=[hash[key]])
: +- TableSourceScan(table=[[test-catalog, default, src2, project=[key]]],
fields=[key])
+- Exchange(distribution=[hash[key]])
+- TableSourceScan(table=[[test-catalog, default, src1]], fields=[key,
val])
{noformat}
Stack trace is:
{noformat}
Caused by: java.lang.ClassCastException: org.apache.flink.types.Row cannot be
cast to org.apache.flink.table.data.RowData
at
org.apache.flink.streaming.api.operators.StreamMap.processElement(StreamMap.java:41)
at
org.apache.flink.table.runtime.operators.multipleinput.output.OneInputStreamOperatorOutput.pushToOperator(OneInputStreamOperatorOutput.java:77)
at
org.apache.flink.table.runtime.operators.multipleinput.output.OneInputStreamOperatorOutput.collect(OneInputStreamOperatorOutput.java:62)
at
org.apache.flink.table.runtime.operators.multipleinput.output.OneInputStreamOperatorOutput.collect(OneInputStreamOperatorOutput.java:32)
at
org.apache.flink.table.runtime.operators.multipleinput.output.BroadcastingOutput.collect(BroadcastingOutput.java:74)
at
org.apache.flink.table.runtime.operators.multipleinput.output.BroadcastingOutput.collect(BroadcastingOutput.java:37)
at
org.apache.flink.streaming.api.operators.CountingOutput.collect(CountingOutput.java:52)
at
org.apache.flink.streaming.api.operators.CountingOutput.collect(CountingOutput.java:30)
at BatchExecCalc$16.processElement(Unknown Source)
at
org.apache.flink.table.runtime.operators.multipleinput.output.OneInputStreamOperatorOutput.pushToOperator(OneInputStreamOperatorOutput.java:77)
at
org.apache.flink.table.runtime.operators.multipleinput.output.OneInputStreamOperatorOutput.collect(OneInputStreamOperatorOutput.java:62)
at
org.apache.flink.table.runtime.operators.multipleinput.output.OneInputStreamOperatorOutput.collect(OneInputStreamOperatorOutput.java:32)
at
org.apache.flink.streaming.api.operators.CountingOutput.collect(CountingOutput.java:52)
at
org.apache.flink.streaming.api.operators.CountingOutput.collect(CountingOutput.java:30)
at
org.apache.flink.table.runtime.util.StreamRecordCollector.collect(StreamRecordCollector.java:44)
at
org.apache.flink.table.runtime.operators.join.HashJoinOperator.collect(HashJoinOperator.java:201)
at
org.apache.flink.table.runtime.operators.join.HashJoinOperator.innerJoin(HashJoinOperator.java:184)
at
org.apache.flink.table.runtime.operators.join.HashJoinOperator$BuildOuterHashJoinOperator.join(HashJoinOperator.java:331)
at
org.apache.flink.table.runtime.operators.join.HashJoinOperator.joinWithNextKey(HashJoinOperator.java:178)
at
org.apache.flink.table.runtime.operators.join.HashJoinOperator.processElement2(HashJoinOperator.java:147)
at
org.apache.flink.table.runtime.operators.multipleinput.input.SecondInputOfTwoInput.processElement(SecondInputOfTwoInput.java:41)
at
org.apache.flink.streaming.runtime.io.StreamMultipleInputProcessorFactory$StreamTaskNetworkOutput.emitRecord(StreamMultipleInputProcessorFactory.java:259)
at
org.apache.flink.streaming.runtime.io.StreamTaskNetworkInput.processElement(StreamTaskNetworkInput.java:184)
at
org.apache.flink.streaming.runtime.io.StreamTaskNetworkInput.emitNext(StreamTaskNetworkInput.java:157)
at
org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:67)
at
org.apache.flink.streaming.runtime.io.StreamMultipleInputProcessor.processInput(StreamMultipleInputProcessor.java:83)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:372)
at
org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:191)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:575)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:539)
at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:722)
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:547)
at java.lang.Thread.run(Thread.java:748)
{noformat}
was (Author: lirui):
Stack trace is:
{noformat}
Caused by: java.lang.ClassCastException: org.apache.flink.types.Row cannot be
cast to org.apache.flink.table.data.RowData
at
org.apache.flink.streaming.api.operators.StreamMap.processElement(StreamMap.java:41)
at
org.apache.flink.table.runtime.operators.multipleinput.output.OneInputStreamOperatorOutput.pushToOperator(OneInputStreamOperatorOutput.java:77)
at
org.apache.flink.table.runtime.operators.multipleinput.output.OneInputStreamOperatorOutput.collect(OneInputStreamOperatorOutput.java:62)
at
org.apache.flink.table.runtime.operators.multipleinput.output.OneInputStreamOperatorOutput.collect(OneInputStreamOperatorOutput.java:32)
at
org.apache.flink.table.runtime.operators.multipleinput.output.BroadcastingOutput.collect(BroadcastingOutput.java:74)
at
org.apache.flink.table.runtime.operators.multipleinput.output.BroadcastingOutput.collect(BroadcastingOutput.java:37)
at
org.apache.flink.streaming.api.operators.CountingOutput.collect(CountingOutput.java:52)
at
org.apache.flink.streaming.api.operators.CountingOutput.collect(CountingOutput.java:30)
at BatchExecCalc$16.processElement(Unknown Source)
at
org.apache.flink.table.runtime.operators.multipleinput.output.OneInputStreamOperatorOutput.pushToOperator(OneInputStreamOperatorOutput.java:77)
at
org.apache.flink.table.runtime.operators.multipleinput.output.OneInputStreamOperatorOutput.collect(OneInputStreamOperatorOutput.java:62)
at
org.apache.flink.table.runtime.operators.multipleinput.output.OneInputStreamOperatorOutput.collect(OneInputStreamOperatorOutput.java:32)
at
org.apache.flink.streaming.api.operators.CountingOutput.collect(CountingOutput.java:52)
at
org.apache.flink.streaming.api.operators.CountingOutput.collect(CountingOutput.java:30)
at
org.apache.flink.table.runtime.util.StreamRecordCollector.collect(StreamRecordCollector.java:44)
at
org.apache.flink.table.runtime.operators.join.HashJoinOperator.collect(HashJoinOperator.java:201)
at
org.apache.flink.table.runtime.operators.join.HashJoinOperator.innerJoin(HashJoinOperator.java:184)
at
org.apache.flink.table.runtime.operators.join.HashJoinOperator$BuildOuterHashJoinOperator.join(HashJoinOperator.java:331)
at
org.apache.flink.table.runtime.operators.join.HashJoinOperator.joinWithNextKey(HashJoinOperator.java:178)
at
org.apache.flink.table.runtime.operators.join.HashJoinOperator.processElement2(HashJoinOperator.java:147)
at
org.apache.flink.table.runtime.operators.multipleinput.input.SecondInputOfTwoInput.processElement(SecondInputOfTwoInput.java:41)
at
org.apache.flink.streaming.runtime.io.StreamMultipleInputProcessorFactory$StreamTaskNetworkOutput.emitRecord(StreamMultipleInputProcessorFactory.java:259)
at
org.apache.flink.streaming.runtime.io.StreamTaskNetworkInput.processElement(StreamTaskNetworkInput.java:184)
at
org.apache.flink.streaming.runtime.io.StreamTaskNetworkInput.emitNext(StreamTaskNetworkInput.java:157)
at
org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:67)
at
org.apache.flink.streaming.runtime.io.StreamMultipleInputProcessor.processInput(StreamMultipleInputProcessor.java:83)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:372)
at
org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:191)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:575)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:539)
at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:722)
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:547)
at java.lang.Thread.run(Thread.java:748)
{noformat}
> Fail to insert into hive table due to ClassCastException
> --------------------------------------------------------
>
> Key: FLINK-20722
> URL: https://issues.apache.org/jira/browse/FLINK-20722
> Project: Flink
> Issue Type: Bug
> Components: Table SQL / Runtime
> Reporter: Rui Li
> Priority: Major
> Fix For: 1.13.0
>
>
> Add the following test in {{TableEnvHiveConnectorITCase}} to reproduce the
> issue:
> {code}
> @Test
> public void test() throws Exception {
> TableEnvironment tableEnv = getTableEnvWithHiveCatalog();
> tableEnv.executeSql("create table src1(key string, val
> string)");
> tableEnv.executeSql("create table src2(key string, val
> string)");
> tableEnv.executeSql("create table dest(key string, val
> string)");
> HiveTestUtils.createTextTableInserter(hiveCatalog, "default",
> "src1")
> .addRow(new Object[]{"1", "val1"})
> .addRow(new Object[]{"2", "val2"})
> .addRow(new Object[]{"3", "val3"})
> .commit();
> HiveTestUtils.createTextTableInserter(hiveCatalog, "default",
> "src2")
> .addRow(new Object[]{"3", "val4"})
> .addRow(new Object[]{"4", "val4"})
> .commit();
> tableEnv.executeSql("INSERT OVERWRITE dest\n" +
> "SELECT j.*\n" +
> "FROM (SELECT t1.key, p1.val\n" +
> " FROM src2 t1\n" +
> " LEFT OUTER JOIN src1 p1\n" +
> " ON (t1.key = p1.key)\n" +
> " UNION ALL\n" +
> " SELECT t2.key, p2.val\n" +
> " FROM src2 t2\n" +
> " LEFT OUTER JOIN src1 p2\n" +
> " ON (t2.key = p2.key)) j").await();
> }
> {code}
--
This message was sent by Atlassian Jira
(v8.3.4#803005)