TableAPI SQL left join error "doesn't support consuming update and delete changes which is produced by node Join(joinType=[LeftOuterJoin]..."

Viewed 86

Need to left join 2 data streams and the Table API throws an exception during starting the Flink app.

It is needed to calculate based on 2 different windows on a same input stream. So eventually the 2 calculated tables/steams need to be joined.

If I use "join" instead of using "left join" the start-up error disappears. Any idea how I can solve it? Is "left join" not supported?

Below is the exception:

org.apache.flink.client.program.ProgramInvocationException: The main method caused an error: Table sink 'default_catalog.default_database.Unregistered_DataStream_Sink_1' doesn't support consuming update and delete changes which is produced by node Join(joinType=[LeftOuterJoin], where=[((userId = userId0) AND (eventEpoch = eventEpoch0))], select=[userId, eventEpoch, AVG_VALUE_A, FIRST_VALUE_B, LAST_VALUE_B, FIRST_EVENTTIME, LAST_EVENTTIME, userId0, eventEpoch0], leftInputSpec=[NoUniqueKey], rightInputSpec=[NoUniqueKey])\n\tat org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:372)\n\tat org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:222)\n\tat

Below is the code snippet.

Table inputTable = tableEnv.fromDataStream(
    inputStream, 
  Schema.newBuilder()
      .columnByExpression("rowtime", "CAST(eventTime AS TIMESTAMP(3))")
      .watermark("rowtime", "rowtime - INTERVAL '" + WATERMARK_DELAY_IN_SECOND + "' SECOND")
      .build());

//Register the table for use in SQL queries
tableEnv.createTemporaryView("InputTable", inputTable);

Table resultTable = tableEnv.sqlQuery(
    "SELECT " +
    " userId," +
    " eventEpoch," +
    " AVG(valueA) OVER w1 AS AVG_VALUE_B" +
    " FROM InputTable " +
    " WINDOW w1 AS (" +
    " PARTITION BY userId ORDER BY rowtime" +
    " RANGE BETWEEN INTERVAL '" + INTERVAL_A + "' SECOND PRECEDING AND CURRENT ROW)"
);

Table resultTable2 = tableEnv.sqlQuery(
    "SELECT " +
    " FIRST_VALUE(valueB) OVER w1 AS FIRST_VALUE_B," +
    " LAST_VALUE(valueB) OVER w1 AS LAST_VALUE_B," +
    " FIRST_VALUE(eventEpoch) OVER w1 AS FIRST_EVENTTIME," +
    " LAST_VALUE(eventEpoch) OVER w1 AS LAST_EVENTTIME," +
    " userId, eventEpoch" +
    " FROM InputTable" +
    " WINDOW w1 AS (" +
    " PARTITION BY userId ORDER BY rowtime" +
    " RANGE BETWEEN INTERVAL '" + INTERVAL_B + "' SECOND PRECEDING AND CURRENT ROW)"
);

tableEnv.createTemporaryView("ResultTable1", resultTable);
tableEnv.createTemporaryView("ResultTable2", resultTable2);

Table resultTable3 = tableEnv.sqlQuery(
    "SELECT ResultTable1.*," +
    " ResultTable2.FIRST_VALUE_B," +
    " ResultTable2.LAST_VALUE_B," +
    " ResultTable2.FIRST_EVENTTIME," +
    " ResultTable2.LAST_EVENTTIME" +
    " FROM ResultTable1" +
    " LEFT JOIN ResultTable2" +
    " ON ResultTable1.userId = ResultTable2.userId" +
    " AND ResultTable1.eventEpoch = ResultTable2.eventEpoch"
);
0 Answers
Related