Skip to content

Commit 551c329

Browse files
committed
[hotfix-32959][core][launcher]fix group by window no retract
1 parent b74452d commit 551c329

File tree

2 files changed

+2
-1
lines changed

2 files changed

+2
-1
lines changed

core/src/main/java/com/dtstack/flink/sql/side/SideSqlExec.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -525,6 +525,7 @@ private void joinFun(Object pollObj,
525525
RowTypeInfo typeInfo = new RowTypeInfo(fieldDataTypes, targetTable.getSchema().getFieldNames());
526526

527527
DataStream<BaseRow> adaptStream = tableEnv.toRetractStream(targetTable, typeInfo)
528+
.filter(f -> f.f0)
528529
.map(f -> RowDataConvert.convertToBaseRow(f));
529530

530531
//join side table before keyby ===> Reducing the size of each dimension table cache of async

launcher/src/main/java/org/apache/flink/table/planner/plan/QueryOperationConverter.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -442,7 +442,7 @@ private RelNode convertToDataStreamScan(
442442
rowType,
443443
dataStream,
444444
false,
445-
true,
445+
false,
446446
fieldIndices,
447447
tableSchema.getFieldNames(),
448448
FlinkStatistic.UNKNOWN(),

0 commit comments

Comments
 (0)