gustavodemorais commented on code in PR #27886:
URL: https://github.com/apache/flink/pull/27886#discussion_r3116885059


##########
flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/SystemTypeInference.java:
##########
@@ -314,13 +336,16 @@ private List<Field> derivePassThroughFields(CallContext 
callContext) {
                                 if 
(arg.is(StaticArgumentTrait.PASS_COLUMNS_THROUGH)) {
                                     return 
DataType.getFields(argDataTypes.get(pos)).stream();
                                 }
-                                if 
(!arg.is(StaticArgumentTrait.SET_SEMANTIC_TABLE)) {
+                                final TableSemantics semantics =
+                                        
callContext.getTableSemantics(pos).orElse(null);
+                                if (semantics == null) {
+                                    return Stream.<Field>empty();
+                                }
+                                final TraitContext traitCtx =

Review Comment:
   Done



##########
flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/SystemTypeInference.java:
##########
@@ -656,12 +684,13 @@ private static void checkRowSemantics(StaticArgument 
staticArg, TableSemantics s
             }
         }
 
-        private static void checkSetSemantics(StaticArgument staticArg, 
TableSemantics semantics) {
-            if (!staticArg.is(StaticArgumentTrait.SET_SEMANTIC_TABLE)) {
+        private static void checkSetSemantics(
+                StaticArgument staticArg, TableSemantics semantics, 
TraitContext traitCtx) {

Review Comment:
   Done



##########
flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/TraitCondition.java:
##########
@@ -0,0 +1,69 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.table.types.inference;
+
+import org.apache.flink.annotation.PublicEvolving;
+
+import java.io.Serializable;
+
+/**
+ * A condition that determines whether a conditional trait on a {@link 
StaticArgument} should be
+ * active for a given call.
+ *
+ * <p>Conditions are evaluated at planning time using the {@link TraitContext} 
which provides access
+ * to the SQL call's properties (PARTITION BY presence, scalar argument 
values, etc.).

Review Comment:
   Done



##########
flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/TraitCondition.java:
##########
@@ -0,0 +1,69 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.table.types.inference;
+
+import org.apache.flink.annotation.PublicEvolving;
+
+import java.io.Serializable;
+
+/**
+ * A condition that determines whether a conditional trait on a {@link 
StaticArgument} should be
+ * active for a given call.
+ *
+ * <p>Conditions are evaluated at planning time using the {@link TraitContext} 
which provides access
+ * to the SQL call's properties (PARTITION BY presence, scalar argument 
values, etc.).
+ *
+ * <p>Use the static factory methods for common conditions:
+ *
+ * <pre>{@code
+ * import static org.apache.flink.table.types.inference.TraitCondition.*;
+ *
+ * StaticArgument.table("input", Row.class, false, EnumSet.of(TABLE, 
SUPPORT_UPDATES))
+ *         .addTraitWhen(hasPartitionBy(), SET_SEMANTIC_TABLE)
+ *         .addTraitWhen(not(hasPartitionBy()), ROW_SEMANTIC_TABLE)
+ *         .addTraitWhen(argIsTrue("produces_full_deletes"), 
REQUIRE_UPDATE_BEFORE);
+ * }</pre>
+ */
+@PublicEvolving
+@FunctionalInterface
+public interface TraitCondition extends Serializable {
+
+    /** Evaluates this condition against the given context. */
+    boolean test(TraitContext ctx);
+
+    /** True when PARTITION BY is provided on the table argument. */
+    static TraitCondition hasPartitionBy() {
+        return TraitContext::hasPartitionBy;
+    }
+
+    /** True when the named boolean argument is provided and its value is 
{@code true}. */
+    static TraitCondition argIsTrue(final String name) {

Review Comment:
   Done



##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecProcessTableFunction.java:
##########
@@ -276,15 +279,17 @@ private RuntimeTableSemantics createRuntimeTableSemantics(
         }
 
         final int timeColumn = 
inputTimeColumns.get(tableArgCall.getInputIndex());
+        final org.apache.flink.table.types.inference.TraitContext traitCtx =

Review Comment:
   Done



##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecProcessTableFunction.java:
##########
@@ -276,15 +279,17 @@ private RuntimeTableSemantics createRuntimeTableSemantics(
         }
 
         final int timeColumn = 
inputTimeColumns.get(tableArgCall.getInputIndex());
+        final org.apache.flink.table.types.inference.TraitContext traitCtx =

Review Comment:
   Done



##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalProcessTableFunction.java:
##########
@@ -221,7 +223,11 @@ protected RelDataType deriveRowType() {
         if (uidRexNode.getKind() == SqlKind.DEFAULT) {
             // Optional for constant or row semantics functions
             if (staticArgs.stream()
-                    .noneMatch(arg -> 
arg.is(StaticArgumentTrait.SET_SEMANTIC_TABLE))) {
+                    .noneMatch(
+                            arg ->
+                                    
arg.is(StaticArgumentTrait.SET_SEMANTIC_TABLE)
+                                            || arg.hasConditionalTrait(

Review Comment:
   Done



##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalProcessTableFunction.java:
##########
@@ -348,6 +354,54 @@ public static List<Ord<StaticArgument>> 
getProvidedInputArgs(RexCall call) {
                 .collect(Collectors.toList());
     }
 
+    /**
+     * Builds a {@link TraitContext} for resolving conditional traits on a 
table argument at
+     * planning time.
+     */
+    public static TraitContext buildTraitContext(
+            final RexCall call, final RexTableArgCall tableArgCall) {
+        final List<StaticArgument> declaredArgs = getStaticArguments(call);
+        final List<RexNode> operands = call.getOperands();
+
+        return new TraitContext() {
+            @Override
+            public boolean hasPartitionBy() {
+                return tableArgCall.getPartitionKeys().length > 0;
+            }
+
+            @Override
+            public <T> Optional<T> getScalarArgument(final String name, final 
Class<T> clazz) {
+                return findScalarLiteral(declaredArgs, operands, name, clazz);
+            }
+        };
+    }
+
+    private static List<StaticArgument> getStaticArguments(final RexCall call) 
{
+        final BridgingSqlFunction.WithTableFunction function =
+                (BridgingSqlFunction.WithTableFunction) call.getOperator();
+        return function.getTypeInference()
+                .getStaticArguments()
+                .orElseThrow(IllegalStateException::new);
+    }
+
+    @SuppressWarnings("unchecked")
+    private static <T> Optional<T> findScalarLiteral(
+            final List<StaticArgument> declaredArgs,
+            final List<RexNode> operands,
+            final String name,
+            final Class<T> clazz) {
+        for (int i = 0; i < declaredArgs.size(); i++) {
+            if (declaredArgs.get(i).getName().equals(name)) {
+                final RexNode operand = operands.get(i);
+                if (operand.getKind() == SqlKind.DEFAULT || !(operand 
instanceof RexLiteral)) {
+                    return Optional.empty();
+                }
+                return Optional.ofNullable(((RexLiteral) 
operand).getValueAs(clazz));

Review Comment:
   Done



##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/rules/physical/stream/StreamPhysicalProcessTableFunctionRule.java:
##########
@@ -127,9 +138,17 @@ private static RelNode applyDistributionOnInput(
     }
 
     private static FlinkRelDistribution deriveDistribution(
-            RexTableArgCall tableOperand, TableCharacteristic 
tableCharacteristic) {
-        if (tableCharacteristic.semantics == Semantics.SET) {
-            final int[] partitionKeys = tableOperand.getPartitionKeys();
+            RexTableArgCall tableOperand,
+            TableCharacteristic tableCharacteristic,
+            StaticArgument staticArg) {
+        final int[] partitionKeys = tableOperand.getPartitionKeys();
+        final boolean hasPartitionBy = partitionKeys.length > 0;
+        final boolean reportedAsSet = tableCharacteristic.semantics == 
Semantics.SET;
+        final boolean setIsConditional =
+                
staticArg.hasConditionalTrait(StaticArgumentTrait.SET_SEMANTIC_TABLE);

Review Comment:
   Done



##########
flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/StaticArgument.java:
##########
@@ -57,18 +61,58 @@ public class StaticArgument {
     private final @Nullable Class<?> conversionClass;
     private final boolean isOptional;
     private final EnumSet<StaticArgumentTrait> traits;
+    private final List<ConditionalTrait> conditionalTraits;
+
+    /** A trait that is conditionally added based on a {@link TraitCondition}. 
*/
+    private static final class ConditionalTrait implements Serializable {
+        private final TraitCondition condition;
+        private final StaticArgumentTrait trait;
+
+        ConditionalTrait(final TraitCondition condition, final 
StaticArgumentTrait trait) {
+            this.condition = condition;
+            this.trait = trait;
+        }
+
+        @Override
+        public boolean equals(final Object o) {
+            if (this == o) {
+                return true;
+            }
+            if (o == null || getClass() != o.getClass()) {
+                return false;
+            }
+            final ConditionalTrait that = (ConditionalTrait) o;
+            return Objects.equals(condition, that.condition) && trait == 
that.trait;
+        }
+
+        @Override
+        public int hashCode() {
+            return Objects.hash(condition, trait);
+        }
+    }
 
     private StaticArgument(
             String name,
             @Nullable DataType dataType,
             @Nullable Class<?> conversionClass,
             boolean isOptional,
             EnumSet<StaticArgumentTrait> traits) {
+        this(name, dataType, conversionClass, isOptional, traits, List.of());
+    }
+
+    private StaticArgument(
+            String name,
+            @Nullable DataType dataType,
+            @Nullable Class<?> conversionClass,
+            boolean isOptional,
+            EnumSet<StaticArgumentTrait> traits,
+            List<ConditionalTrait> conditionalTraits) {
         this.name = Preconditions.checkNotNull(name, "Name must not be null.");
         this.dataType = dataType;
         this.conversionClass = conversionClass;
         this.isOptional = isOptional;
         this.traits = Preconditions.checkNotNull(traits, "Traits must not be 
null.");
+        this.conditionalTraits = 
Collections.unmodifiableList(conditionalTraits);

Review Comment:
   Done



-- 
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]

Reply via email to