gustavodemorais commented on code in PR #29335:
URL: https://github.com/apache/flink/pull/29335#discussion_r4153679478
##########
flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/SystemTypeInference.java:
##########
@@ -348,6 +348,12 @@ public Optional<DataType> inferType(CallContext
callContext) {
fields.addAll(deriveRowtimeField(callContext, resolvedArgs));
}
+ if (fields.isEmpty()) {
+ // Only a fully empty row falls back to
EXPR$0, for backwards
+ // compatibility.
+ fields.add(DataTypes.FIELD("EXPR$0",
functionDataType));
+ }
+
Review Comment:
```suggestion
// Nothing to show at all: no pass-through,
no rowtime, and the
// function itself returned ROW<>. Fall back
to EXPR$0 as output.
if (fields.isEmpty()) {
fields.add(DataTypes.FIELD("EXPR$0",
functionDataType));
}
```
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java:
##########
@@ -791,6 +792,29 @@ public class ProcessTableFunctionTestPrograms {
.runSql("INSERT INTO sink SELECT * FROM f(r => TABLE t
PARTITION BY name)")
.build();
+ public static final TableTestProgram PROCESS_EMPTY_OUTPUT =
Review Comment:
Could you add a rowtime-only variant? Add a sibling program right after
PROCESS_EMPTY_OUTPUT, dropping PARTITION BY
##########
flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/SystemTypeInference.java:
##########
@@ -403,16 +409,11 @@ private List<Field> derivePassThroughFields(
}
private List<Field> deriveFunctionOutputFields(DataType
functionDataType) {
- final List<DataType> fieldTypes =
DataType.getFieldDataTypes(functionDataType);
- final List<String> fieldNames =
DataType.getFieldNames(functionDataType);
-
- if (fieldTypes.isEmpty()) {
- // Before the system type inference was introduced, SQL and
- // Table API chose a different default field name.
- // EXPR$0 is chosen for best-effort backwards compatibility for
- // SQL users.
+ if
(!LogicalTypeChecks.isCompositeType(functionDataType.getLogicalType())) {
return List.of(DataTypes.FIELD("EXPR$0", functionDataType));
}
+ final List<DataType> fieldTypes =
DataType.getFieldDataTypes(functionDataType);
+ final List<String> fieldNames =
DataType.getFieldNames(functionDataType);
return IntStream.range(0, fieldTypes.size())
.mapToObj(pos -> DataTypes.FIELD(fieldNames.get(pos),
fieldTypes.get(pos)))
.collect(Collectors.toList());
Review Comment:
```suggestion
final boolean isScalarOutput =
!LogicalTypeChecks.isCompositeType(functionDataType.getLogicalType());
if (isScalarOutput) {
// Scalar output has no field name of its own; EXPR$0 is kept for
// backwards compatibility with the pre-system-type-inference
default.
return List.of(DataTypes.FIELD("EXPR$0", functionDataType));
}
// For composite types, extract field names/types and build output
final List<DataType> fieldTypes =
DataType.getFieldDataTypes(functionDataType);
final List<String> fieldNames = DataType.getFieldNames(functionDataType);
return IntStream.range(0, fieldTypes.size())
.mapToObj(pos -> DataTypes.FIELD(fieldNames.get(pos),
fieldTypes.get(pos)))
.collect(Collectors.toList());
}
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]