[ 
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)

Reply via email to