suxiaogang223 commented on code in PR #68320:
URL: https://github.com/apache/doris/pull/68320#discussion_r4090757618
##########
be/src/exec/operator/exchange_sink_operator.cpp:
##########
@@ -535,9 +567,11 @@ Status ExchangeSinkOperatorX::sink_impl(RuntimeState*
state, Block* block, bool
_part_type == TPartitionType::BUCKET_SHFFULE_HASH_PARTITIONED ||
_part_type == TPartitionType::OLAP_TABLE_SINK_HASH_PARTITIONED
||
_part_type == TPartitionType::HIVE_TABLE_SINK_HASH_PARTITIONED
||
+ _part_type ==
TPartitionType::EXTERNAL_TABLE_SINK_HASH_PARTITIONED ||
_part_type == TPartitionType::MERGE_PARTITIONED) {
RETURN_IF_ERROR(local_state._writer->write(state, block, eos));
- } else if (_part_type == TPartitionType::HIVE_TABLE_SINK_UNPARTITIONED) {
+ } else if (_part_type == TPartitionType::HIVE_TABLE_SINK_UNPARTITIONED ||
+ _part_type ==
TPartitionType::EXTERNAL_TABLE_SINK_UNPARTITIONED) {
Review Comment:
Addressed in eff14993d69. External unpartitioned writer scaling now updates
shared byte/count state under AnnotatedMutex, and ExchangeSinkOperatorTest
covers concurrent local-state updates.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/plugin/PluginDrivenExternalTable.java:
##########
@@ -423,6 +507,29 @@ public boolean requirePartitionHashOnWrite() {
.orElse(false);
}
+ /** Returns this table's connector-owned write distribution, or empty for
generic planning. */
+ public Optional<ConnectorWriteDistribution>
getConnectorWriteDistribution() {
+ if (!(catalog instanceof PluginDrivenExternalCatalog)) {
+ return Optional.empty();
+ }
+ PluginDrivenExternalCatalog pluginCatalog =
(PluginDrivenExternalCatalog) catalog;
+ Connector connector = pluginCatalog.getConnector();
+ if (connector == null) {
+ return Optional.empty();
+ }
+ ConnectorSession session = pluginCatalog.buildConnectorSession();
+ ConnectorMetadata metadata = PluginDrivenMetadata.get(session,
connector);
+ Optional<ConnectorTableHandle> handle =
resolveConnectorTableHandle(session, metadata);
+ if (!handle.isPresent()) {
+ return Optional.empty();
+ }
+ ConnectorWritePlanProvider provider =
connector.getWritePlanProvider(handle.get());
+ if (provider == null) {
+ return Optional.empty();
+ }
+ return Optional.ofNullable(provider.getWriteDistribution(session,
handle.get()));
Review Comment:
Addressed in 70d3c439123. Provider resolution and connector write callbacks
now use one restoring plugin-TCCL helper, with tests asserting the observed
loader and restoration.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/ConnectorChangelogPlanBuilder.java:
##########
@@ -0,0 +1,547 @@
+// 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.doris.nereids.rules.analysis;
+
+import org.apache.doris.catalog.Column;
+import org.apache.doris.connector.spi.write.ConnectorChangelogMode;
+import org.apache.doris.nereids.CascadesContext;
+import org.apache.doris.nereids.analyzer.Scope;
+import org.apache.doris.nereids.analyzer.UnboundAlias;
+import org.apache.doris.nereids.analyzer.UnboundSlot;
+import org.apache.doris.nereids.exceptions.AnalysisException;
+import org.apache.doris.nereids.trees.expressions.Alias;
+import org.apache.doris.nereids.trees.expressions.Cast;
+import org.apache.doris.nereids.trees.expressions.EqualTo;
+import org.apache.doris.nereids.trees.expressions.Expression;
+import org.apache.doris.nereids.trees.expressions.LessThanEqual;
+import org.apache.doris.nereids.trees.expressions.NamedExpression;
+import org.apache.doris.nereids.trees.expressions.Not;
+import org.apache.doris.nereids.trees.expressions.Slot;
+import org.apache.doris.nereids.trees.expressions.WindowExpression;
+import org.apache.doris.nereids.trees.expressions.functions.agg.AnyValue;
+import org.apache.doris.nereids.trees.expressions.functions.agg.Count;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.AssertTrue;
+import
org.apache.doris.nereids.trees.expressions.functions.scalar.ShortCircuitIf;
+import org.apache.doris.nereids.trees.expressions.literal.BigIntLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.IntegerLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.NullLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.TinyIntLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.VarcharLiteral;
+import
org.apache.doris.nereids.trees.plans.commands.info.ConnectorChangelogRowChangeSpec;
+import org.apache.doris.nereids.trees.plans.commands.merge.MergeMatchedClause;
+import
org.apache.doris.nereids.trees.plans.commands.merge.MergeNotMatchedClause;
+import org.apache.doris.nereids.trees.plans.logical.LogicalAggregate;
+import org.apache.doris.nereids.trees.plans.logical.LogicalFilter;
+import org.apache.doris.nereids.trees.plans.logical.LogicalJoin;
+import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
+import org.apache.doris.nereids.trees.plans.logical.LogicalProject;
+import org.apache.doris.nereids.trees.plans.logical.LogicalWindow;
+import org.apache.doris.nereids.types.BigIntType;
+import org.apache.doris.nereids.types.DataType;
+import org.apache.doris.nereids.types.IntegerType;
+import org.apache.doris.nereids.util.ExpressionUtils;
+import org.apache.doris.nereids.util.TypeCoercionUtils;
+
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableSet;
+import com.google.common.collect.Lists;
+import com.google.common.collect.Maps;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.TreeSet;
+
+/** Builds the operation-column plus full-row projection used by
changelog-oriented connectors. */
+public final class ConnectorChangelogPlanBuilder {
+ private static final String BRANCH_LABEL = "__DORIS_CHANGELOG_BRANCH__";
+
+ private ConnectorChangelogPlanBuilder() {
+ }
+
+ /** Builds a changelog plan for the requested connector row-level
operation. */
+ public static LogicalPlan build(List<Column> schema, List<String>
primaryKeys,
+ ConnectorChangelogMode mode, ConnectorChangelogRowChangeSpec spec,
+ LogicalPlan child, CascadesContext context) {
+ if (spec instanceof ConnectorChangelogRowChangeSpec.Update) {
+ return buildUpdate(schema, mode,
(ConnectorChangelogRowChangeSpec.Update) spec,
+ child, context);
+ }
+ if (spec instanceof ConnectorChangelogRowChangeSpec.Delete) {
+ return buildDelete(schema, primaryKeys, mode,
+ (ConnectorChangelogRowChangeSpec.Delete) spec, child,
context);
+ }
+ if (spec instanceof ConnectorChangelogRowChangeSpec.Merge) {
+ return new MergeBuilder(schema, primaryKeys, mode,
+ (ConnectorChangelogRowChangeSpec.Merge) spec, child,
context).build();
+ }
+ throw new AnalysisException("Unsupported connector changelog
specification: "
+ + spec.getClass().getSimpleName());
+ }
+
+ private static LogicalPlan buildUpdate(List<Column> schema,
ConnectorChangelogMode mode,
+ ConnectorChangelogRowChangeSpec.Update update, LogicalPlan child,
+ CascadesContext context) {
+ Map<String, Expression> changes =
Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER);
+ for (EqualTo assignment : update.getAssignments()) {
+ List<String> parts = ((UnboundSlot)
assignment.left()).getNameParts();
+ String name = parts.get(parts.size() - 1);
+ if (changes.put(name, assignment.right()) != null) {
+ throw new AnalysisException("Duplicate column name in
connector UPDATE: " + name);
+ }
+ }
+ ExpressionAnalyzer analyzer = analyzer(child, context);
+ List<NamedExpression> projects = new ArrayList<>();
+ projects.add(operation(mode.getOperationColumnName(),
mode.getUpdateValue()));
+ for (Column column : schema) {
+ Expression value = changes.remove(column.getName());
+ if (value == null) {
+ value = targetSlot(update.getTargetNameInPlan(),
column.getName());
+ }
+ projects.add(bindColumn(analyzer, value, column));
+ }
+ if (!changes.isEmpty()) {
+ throw new AnalysisException("Unknown column in connector UPDATE: "
+ + String.join(", ", changes.keySet()));
+ }
+ return new LogicalProject<>(projects, child);
+ }
+
+ private static LogicalPlan buildDelete(List<Column> schema, List<String>
primaryKeys,
+ ConnectorChangelogMode mode,
+ ConnectorChangelogRowChangeSpec.Delete delete, LogicalPlan child,
+ CascadesContext context) {
+ ExpressionAnalyzer analyzer = analyzer(child, context);
+ List<NamedExpression> projects = new ArrayList<>();
+ projects.add(operation(mode.getOperationColumnName(),
mode.getDeleteValue()));
+ for (Column column : schema) {
+ projects.add(bindColumn(analyzer,
+ targetSlot(delete.getTargetNameInPlan(),
column.getName()), column));
+ }
+ LogicalProject<LogicalPlan> project = new LogicalProject<>(projects,
child);
+ if (!delete.shouldDeduplicateTargetRows()) {
+ return project;
+ }
+ if (primaryKeys.isEmpty()) {
+ throw new AnalysisException("Connector DELETE USING requires a
primary-key table");
+ }
+ Set<String> keys = new TreeSet<>(String.CASE_INSENSITIVE_ORDER);
+ keys.addAll(primaryKeys);
+ List<Expression> groupBy = new ArrayList<>();
+ List<NamedExpression> outputs = new ArrayList<>();
+ Slot operation = project.getOutput().get(0);
+ groupBy.add(operation);
+ outputs.add(operation);
+ for (int i = 0; i < schema.size(); i++) {
+ Column column = schema.get(i);
+ Slot value = project.getOutput().get(i + 1);
+ if (keys.remove(column.getName())) {
+ groupBy.add(value);
+ outputs.add(value);
+ } else {
+ outputs.add(new Alias(new AnyValue(value), column.getName()));
+ }
+ }
+ if (!keys.isEmpty()) {
+ throw new AnalysisException("Unknown connector primary-key column:
"
+ + String.join(", ", keys));
+ }
+ return new LogicalAggregate<>(groupBy, outputs, project);
+ }
+
+ private static Alias operation(String columnName, byte value) {
+ return new Alias(new TinyIntLiteral(value), columnName);
+ }
+
+ private static UnboundSlot targetSlot(List<String> qualifier, String
column) {
+ List<String> parts = new ArrayList<>(qualifier);
+ parts.add(column);
+ return new UnboundSlot(parts);
+ }
+
+ private static ExpressionAnalyzer analyzer(LogicalPlan plan,
CascadesContext context) {
+ return new ExpressionAnalyzer(plan, new Scope(plan.getOutput()),
context, true, false);
+ }
+
+ private static Alias bindColumn(ExpressionAnalyzer analyzer, Expression
expression, Column column) {
+ Expression value = analyzer.analyze(expression);
+ value = TypeCoercionUtils.castIfNotSameType(value,
DataType.fromCatalogType(column.getType()));
+ return new Alias(value, column.getName());
+ }
+
+ private static final class MergeBuilder {
+ private final List<Column> schema;
+ private final List<String> primaryKeys;
+ private final ConnectorChangelogMode mode;
+ private final ConnectorChangelogRowChangeSpec.Merge merge;
+ private final LogicalPlan child;
+ private final ExpressionAnalyzer analyzer;
+
+ private MergeBuilder(List<Column> schema, List<String> primaryKeys,
+ ConnectorChangelogMode mode,
+ ConnectorChangelogRowChangeSpec.Merge merge, LogicalPlan child,
+ CascadesContext context) {
+ this.schema = schema;
+ this.primaryKeys = primaryKeys;
+ this.mode = mode;
+ this.merge = merge;
+ this.child = child;
+ this.analyzer = analyzer(child, context);
+ }
+
+ private LogicalPlan build() {
+ if (primaryKeys.isEmpty()) {
+ throw new AnalysisException("Connector MERGE requires a
primary-key table");
+ }
+ Alias branch = bindBranchLabel();
+ Slot branchSlot = branch.toSlot();
+ List<NamedExpression> branchOutputs = new
ArrayList<>(child.getOutput());
+ branchOutputs.add(branch);
+ LogicalPlan selected = new LogicalProject<>(branchOutputs, child);
+ selected = new LogicalFilter<>(
+ ImmutableSet.of(new Not(new
org.apache.doris.nereids.trees.expressions.IsNull(branchSlot))),
+ selected);
+ List<List<Expression>> branches = buildBranchProjections();
+ if (!merge.getNotMatchedClauses().isEmpty()) {
+ validateNotMatchedPrimaryKeys(branches);
+ }
+ List<NamedExpression> output = new ArrayList<>();
+ for (int column = 0; column <= schema.size(); column++) {
+ DataType type = column == 0
+ ? org.apache.doris.nereids.types.TinyIntType.INSTANCE
+ : DataType.fromCatalogType(schema.get(column -
1).getType());
+ String name = column == 0
+ ? mode.getOperationColumnName() : schema.get(column -
1).getName();
+ Expression value = new NullLiteral(type);
+ for (int index = branches.size() - 1; index >= 0; index--) {
+ Expression branchValue =
TypeCoercionUtils.castIfNotSameType(
+ branches.get(index).get(column), type);
+ value = new ShortCircuitIf(new EqualTo(branchSlot, new
IntegerLiteral(index)),
+ branchValue, value);
+ }
+ output.add(new Alias(value, name));
+ }
+ return addCardinalityChecks(new LogicalProject<>(output,
selected));
+ }
+
+ private void validateNotMatchedPrimaryKeys(List<List<Expression>>
branches) {
+ Map<String, Slot> targetKeys =
Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER);
+ for (String key : primaryKeys) {
+ targetKeys.put(key, findTargetSlot(key));
+ }
+ Set<Slot> targetSlots = child.getOutput().stream()
+ .filter(slot -> qualifierEndsWith(slot.getQualifier(),
merge.getTargetNameInPlan()))
+ .collect(ImmutableSet.toImmutableSet());
+ if (!(child instanceof LogicalJoin)) {
+ throw new AnalysisException("Connector MERGE input must be a
logical join");
+ }
+ Expression onClause = ((LogicalJoin<?, ?>)
child).getOnClauseCondition()
+ .orElseThrow(() -> new AnalysisException("Connector MERGE
requires an ON condition"));
+ Map<String, Expression> sourceKeys =
Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER);
+ for (Expression conjunct :
ExpressionUtils.extractConjunction(onClause)) {
+ if (!(conjunct instanceof EqualTo)) {
+ throw invalidNotMatchedKeyCondition();
+ }
+ EqualTo equality = (EqualTo) conjunct;
+ String leftKey = targetPrimaryKeyName(equality.left(),
targetKeys);
+ String rightKey = targetPrimaryKeyName(equality.right(),
targetKeys);
+ if ((leftKey == null) == (rightKey == null)) {
+ throw invalidNotMatchedKeyCondition();
+ }
+ String key = leftKey != null ? leftKey : rightKey;
+ Expression source = leftKey != null ? equality.right() :
equality.left();
+ if (source.getInputSlots().isEmpty()
+ ||
source.getInputSlots().stream().anyMatch(targetSlots::contains)
+ || source.containsNondeterministic()
+ || sourceKeys.put(key, source) != null) {
+ throw invalidNotMatchedKeyCondition();
+ }
+ }
+ if (sourceKeys.size() != targetKeys.size()) {
+ throw invalidNotMatchedKeyCondition();
+ }
+ int firstInsert = merge.getMatchedClauses().size();
+ for (int branch = firstInsert; branch < branches.size(); branch++)
{
+ for (Map.Entry<String, Expression> sourceKey :
sourceKeys.entrySet()) {
+ int column = schemaIndex(sourceKey.getKey()) + 1;
+ DataType type = DataType.fromCatalogType(schema.get(column
- 1).getType());
+ if
(!TypeCoercionUtils.castIfNotSameType(branches.get(branch).get(column), type)
+
.equals(TypeCoercionUtils.castIfNotSameType(sourceKey.getValue(), type))) {
+ throw invalidNotMatchedKeyCondition();
+ }
+ }
+ }
+ }
+
+ private LogicalPlan addCardinalityChecks(LogicalProject<?> rowChanges)
{
+ List<Slot> outputs = rowChanges.getOutput();
+ Slot operation = outputs.get(0);
+ List<Expression> partitionKeys = new ArrayList<>();
+ for (String key : primaryKeys) {
+ partitionKeys.add(outputs.get(schemaIndex(key) + 1));
Review Comment:
Addressed in fece43753bd. Matched-row cardinality is partitioned by
preserved original target primary-key aliases; insert cardinality continues to
use emitted keys.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/ConnectorChangelogPlanBuilder.java:
##########
@@ -0,0 +1,547 @@
+// 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.doris.nereids.rules.analysis;
+
+import org.apache.doris.catalog.Column;
+import org.apache.doris.connector.spi.write.ConnectorChangelogMode;
+import org.apache.doris.nereids.CascadesContext;
+import org.apache.doris.nereids.analyzer.Scope;
+import org.apache.doris.nereids.analyzer.UnboundAlias;
+import org.apache.doris.nereids.analyzer.UnboundSlot;
+import org.apache.doris.nereids.exceptions.AnalysisException;
+import org.apache.doris.nereids.trees.expressions.Alias;
+import org.apache.doris.nereids.trees.expressions.Cast;
+import org.apache.doris.nereids.trees.expressions.EqualTo;
+import org.apache.doris.nereids.trees.expressions.Expression;
+import org.apache.doris.nereids.trees.expressions.LessThanEqual;
+import org.apache.doris.nereids.trees.expressions.NamedExpression;
+import org.apache.doris.nereids.trees.expressions.Not;
+import org.apache.doris.nereids.trees.expressions.Slot;
+import org.apache.doris.nereids.trees.expressions.WindowExpression;
+import org.apache.doris.nereids.trees.expressions.functions.agg.AnyValue;
+import org.apache.doris.nereids.trees.expressions.functions.agg.Count;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.AssertTrue;
+import
org.apache.doris.nereids.trees.expressions.functions.scalar.ShortCircuitIf;
+import org.apache.doris.nereids.trees.expressions.literal.BigIntLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.IntegerLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.NullLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.TinyIntLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.VarcharLiteral;
+import
org.apache.doris.nereids.trees.plans.commands.info.ConnectorChangelogRowChangeSpec;
+import org.apache.doris.nereids.trees.plans.commands.merge.MergeMatchedClause;
+import
org.apache.doris.nereids.trees.plans.commands.merge.MergeNotMatchedClause;
+import org.apache.doris.nereids.trees.plans.logical.LogicalAggregate;
+import org.apache.doris.nereids.trees.plans.logical.LogicalFilter;
+import org.apache.doris.nereids.trees.plans.logical.LogicalJoin;
+import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
+import org.apache.doris.nereids.trees.plans.logical.LogicalProject;
+import org.apache.doris.nereids.trees.plans.logical.LogicalWindow;
+import org.apache.doris.nereids.types.BigIntType;
+import org.apache.doris.nereids.types.DataType;
+import org.apache.doris.nereids.types.IntegerType;
+import org.apache.doris.nereids.util.ExpressionUtils;
+import org.apache.doris.nereids.util.TypeCoercionUtils;
+
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableSet;
+import com.google.common.collect.Lists;
+import com.google.common.collect.Maps;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.TreeSet;
+
+/** Builds the operation-column plus full-row projection used by
changelog-oriented connectors. */
+public final class ConnectorChangelogPlanBuilder {
+ private static final String BRANCH_LABEL = "__DORIS_CHANGELOG_BRANCH__";
+
+ private ConnectorChangelogPlanBuilder() {
+ }
+
+ /** Builds a changelog plan for the requested connector row-level
operation. */
+ public static LogicalPlan build(List<Column> schema, List<String>
primaryKeys,
+ ConnectorChangelogMode mode, ConnectorChangelogRowChangeSpec spec,
+ LogicalPlan child, CascadesContext context) {
+ if (spec instanceof ConnectorChangelogRowChangeSpec.Update) {
+ return buildUpdate(schema, mode,
(ConnectorChangelogRowChangeSpec.Update) spec,
+ child, context);
+ }
+ if (spec instanceof ConnectorChangelogRowChangeSpec.Delete) {
+ return buildDelete(schema, primaryKeys, mode,
+ (ConnectorChangelogRowChangeSpec.Delete) spec, child,
context);
+ }
+ if (spec instanceof ConnectorChangelogRowChangeSpec.Merge) {
+ return new MergeBuilder(schema, primaryKeys, mode,
+ (ConnectorChangelogRowChangeSpec.Merge) spec, child,
context).build();
+ }
+ throw new AnalysisException("Unsupported connector changelog
specification: "
+ + spec.getClass().getSimpleName());
+ }
+
+ private static LogicalPlan buildUpdate(List<Column> schema,
ConnectorChangelogMode mode,
+ ConnectorChangelogRowChangeSpec.Update update, LogicalPlan child,
+ CascadesContext context) {
+ Map<String, Expression> changes =
Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER);
+ for (EqualTo assignment : update.getAssignments()) {
+ List<String> parts = ((UnboundSlot)
assignment.left()).getNameParts();
+ String name = parts.get(parts.size() - 1);
+ if (changes.put(name, assignment.right()) != null) {
+ throw new AnalysisException("Duplicate column name in
connector UPDATE: " + name);
+ }
+ }
+ ExpressionAnalyzer analyzer = analyzer(child, context);
+ List<NamedExpression> projects = new ArrayList<>();
+ projects.add(operation(mode.getOperationColumnName(),
mode.getUpdateValue()));
+ for (Column column : schema) {
+ Expression value = changes.remove(column.getName());
+ if (value == null) {
+ value = targetSlot(update.getTargetNameInPlan(),
column.getName());
+ }
+ projects.add(bindColumn(analyzer, value, column));
+ }
+ if (!changes.isEmpty()) {
+ throw new AnalysisException("Unknown column in connector UPDATE: "
+ + String.join(", ", changes.keySet()));
+ }
+ return new LogicalProject<>(projects, child);
+ }
+
+ private static LogicalPlan buildDelete(List<Column> schema, List<String>
primaryKeys,
+ ConnectorChangelogMode mode,
+ ConnectorChangelogRowChangeSpec.Delete delete, LogicalPlan child,
+ CascadesContext context) {
+ ExpressionAnalyzer analyzer = analyzer(child, context);
+ List<NamedExpression> projects = new ArrayList<>();
+ projects.add(operation(mode.getOperationColumnName(),
mode.getDeleteValue()));
+ for (Column column : schema) {
+ projects.add(bindColumn(analyzer,
+ targetSlot(delete.getTargetNameInPlan(),
column.getName()), column));
+ }
+ LogicalProject<LogicalPlan> project = new LogicalProject<>(projects,
child);
+ if (!delete.shouldDeduplicateTargetRows()) {
+ return project;
+ }
+ if (primaryKeys.isEmpty()) {
+ throw new AnalysisException("Connector DELETE USING requires a
primary-key table");
+ }
+ Set<String> keys = new TreeSet<>(String.CASE_INSENSITIVE_ORDER);
+ keys.addAll(primaryKeys);
+ List<Expression> groupBy = new ArrayList<>();
+ List<NamedExpression> outputs = new ArrayList<>();
+ Slot operation = project.getOutput().get(0);
+ groupBy.add(operation);
+ outputs.add(operation);
+ for (int i = 0; i < schema.size(); i++) {
+ Column column = schema.get(i);
+ Slot value = project.getOutput().get(i + 1);
+ if (keys.remove(column.getName())) {
+ groupBy.add(value);
+ outputs.add(value);
+ } else {
+ outputs.add(new Alias(new AnyValue(value), column.getName()));
+ }
+ }
+ if (!keys.isEmpty()) {
+ throw new AnalysisException("Unknown connector primary-key column:
"
+ + String.join(", ", keys));
+ }
+ return new LogicalAggregate<>(groupBy, outputs, project);
+ }
+
+ private static Alias operation(String columnName, byte value) {
+ return new Alias(new TinyIntLiteral(value), columnName);
+ }
+
+ private static UnboundSlot targetSlot(List<String> qualifier, String
column) {
+ List<String> parts = new ArrayList<>(qualifier);
+ parts.add(column);
+ return new UnboundSlot(parts);
+ }
+
+ private static ExpressionAnalyzer analyzer(LogicalPlan plan,
CascadesContext context) {
+ return new ExpressionAnalyzer(plan, new Scope(plan.getOutput()),
context, true, false);
+ }
+
+ private static Alias bindColumn(ExpressionAnalyzer analyzer, Expression
expression, Column column) {
+ Expression value = analyzer.analyze(expression);
+ value = TypeCoercionUtils.castIfNotSameType(value,
DataType.fromCatalogType(column.getType()));
+ return new Alias(value, column.getName());
+ }
+
+ private static final class MergeBuilder {
+ private final List<Column> schema;
+ private final List<String> primaryKeys;
+ private final ConnectorChangelogMode mode;
+ private final ConnectorChangelogRowChangeSpec.Merge merge;
+ private final LogicalPlan child;
+ private final ExpressionAnalyzer analyzer;
+
+ private MergeBuilder(List<Column> schema, List<String> primaryKeys,
+ ConnectorChangelogMode mode,
+ ConnectorChangelogRowChangeSpec.Merge merge, LogicalPlan child,
+ CascadesContext context) {
+ this.schema = schema;
+ this.primaryKeys = primaryKeys;
+ this.mode = mode;
+ this.merge = merge;
+ this.child = child;
+ this.analyzer = analyzer(child, context);
+ }
+
+ private LogicalPlan build() {
+ if (primaryKeys.isEmpty()) {
+ throw new AnalysisException("Connector MERGE requires a
primary-key table");
+ }
+ Alias branch = bindBranchLabel();
+ Slot branchSlot = branch.toSlot();
+ List<NamedExpression> branchOutputs = new
ArrayList<>(child.getOutput());
+ branchOutputs.add(branch);
+ LogicalPlan selected = new LogicalProject<>(branchOutputs, child);
+ selected = new LogicalFilter<>(
+ ImmutableSet.of(new Not(new
org.apache.doris.nereids.trees.expressions.IsNull(branchSlot))),
+ selected);
+ List<List<Expression>> branches = buildBranchProjections();
+ if (!merge.getNotMatchedClauses().isEmpty()) {
+ validateNotMatchedPrimaryKeys(branches);
+ }
+ List<NamedExpression> output = new ArrayList<>();
+ for (int column = 0; column <= schema.size(); column++) {
+ DataType type = column == 0
+ ? org.apache.doris.nereids.types.TinyIntType.INSTANCE
+ : DataType.fromCatalogType(schema.get(column -
1).getType());
+ String name = column == 0
+ ? mode.getOperationColumnName() : schema.get(column -
1).getName();
+ Expression value = new NullLiteral(type);
+ for (int index = branches.size() - 1; index >= 0; index--) {
+ Expression branchValue =
TypeCoercionUtils.castIfNotSameType(
+ branches.get(index).get(column), type);
+ value = new ShortCircuitIf(new EqualTo(branchSlot, new
IntegerLiteral(index)),
+ branchValue, value);
+ }
+ output.add(new Alias(value, name));
+ }
+ return addCardinalityChecks(new LogicalProject<>(output,
selected));
+ }
+
+ private void validateNotMatchedPrimaryKeys(List<List<Expression>>
branches) {
+ Map<String, Slot> targetKeys =
Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER);
+ for (String key : primaryKeys) {
+ targetKeys.put(key, findTargetSlot(key));
+ }
+ Set<Slot> targetSlots = child.getOutput().stream()
+ .filter(slot -> qualifierEndsWith(slot.getQualifier(),
merge.getTargetNameInPlan()))
+ .collect(ImmutableSet.toImmutableSet());
+ if (!(child instanceof LogicalJoin)) {
+ throw new AnalysisException("Connector MERGE input must be a
logical join");
+ }
+ Expression onClause = ((LogicalJoin<?, ?>)
child).getOnClauseCondition()
+ .orElseThrow(() -> new AnalysisException("Connector MERGE
requires an ON condition"));
+ Map<String, Expression> sourceKeys =
Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER);
+ for (Expression conjunct :
ExpressionUtils.extractConjunction(onClause)) {
+ if (!(conjunct instanceof EqualTo)) {
+ throw invalidNotMatchedKeyCondition();
+ }
+ EqualTo equality = (EqualTo) conjunct;
+ String leftKey = targetPrimaryKeyName(equality.left(),
targetKeys);
+ String rightKey = targetPrimaryKeyName(equality.right(),
targetKeys);
+ if ((leftKey == null) == (rightKey == null)) {
+ throw invalidNotMatchedKeyCondition();
+ }
+ String key = leftKey != null ? leftKey : rightKey;
+ Expression source = leftKey != null ? equality.right() :
equality.left();
+ if (source.getInputSlots().isEmpty()
+ ||
source.getInputSlots().stream().anyMatch(targetSlots::contains)
+ || source.containsNondeterministic()
+ || sourceKeys.put(key, source) != null) {
+ throw invalidNotMatchedKeyCondition();
+ }
+ }
+ if (sourceKeys.size() != targetKeys.size()) {
+ throw invalidNotMatchedKeyCondition();
+ }
+ int firstInsert = merge.getMatchedClauses().size();
+ for (int branch = firstInsert; branch < branches.size(); branch++)
{
+ for (Map.Entry<String, Expression> sourceKey :
sourceKeys.entrySet()) {
+ int column = schemaIndex(sourceKey.getKey()) + 1;
+ DataType type = DataType.fromCatalogType(schema.get(column
- 1).getType());
+ if
(!TypeCoercionUtils.castIfNotSameType(branches.get(branch).get(column), type)
+
.equals(TypeCoercionUtils.castIfNotSameType(sourceKey.getValue(), type))) {
+ throw invalidNotMatchedKeyCondition();
+ }
+ }
+ }
+ }
+
+ private LogicalPlan addCardinalityChecks(LogicalProject<?> rowChanges)
{
+ List<Slot> outputs = rowChanges.getOutput();
+ Slot operation = outputs.get(0);
+ List<Expression> partitionKeys = new ArrayList<>();
+ for (String key : primaryKeys) {
+ partitionKeys.add(outputs.get(schemaIndex(key) + 1));
+ }
+ Expression isInsert = new EqualTo(operation, new
TinyIntLiteral(mode.getInsertValue()));
+ List<CardinalityCheck> checks = new ArrayList<>();
+ if (!merge.getMatchedClauses().isEmpty()) {
+ checks.add(CardinalityCheck.matched(isInsert));
+ }
+ if (!merge.getNotMatchedClauses().isEmpty()) {
+ checks.add(CardinalityCheck.inserted(isInsert));
+ }
+ List<NamedExpression> markerOutputs = new ArrayList<>(outputs);
+ for (CardinalityCheck check : checks) {
+ markerOutputs.add(check.marker);
+ }
+ LogicalPlan plan = new LogicalProject<>(markerOutputs, rowChanges);
+ List<Alias> counts = new ArrayList<>();
+ for (CardinalityCheck check : checks) {
+ counts.add(check.count(partitionKeys));
+ }
+ plan = new LogicalWindow<>(new ArrayList<>(counts), plan);
+ ImmutableSet.Builder<Expression> assertions =
ImmutableSet.builder();
+ for (int i = 0; i < checks.size(); i++) {
+ assertions.add(checks.get(i).assertion(counts.get(i)));
+ }
+ plan = new LogicalFilter<>(assertions.build(), plan);
+ return new LogicalProject<>(new ArrayList<>(outputs), plan);
+ }
+
+ private int schemaIndex(String name) {
+ for (int i = 0; i < schema.size(); i++) {
+ if (schema.get(i).getName().equalsIgnoreCase(name)) {
+ return i;
+ }
+ }
+ throw new AnalysisException("Unable to resolve connector
row-change column '" + name + "'");
+ }
+
+ private String targetPrimaryKeyName(Expression expression, Map<String,
Slot> targetKeys) {
+ Expression unwrapped = expression;
+ while (unwrapped instanceof Cast) {
+ if (((Cast) unwrapped).isExplicitType()) {
+ return null;
+ }
+ unwrapped = unwrapped.child(0);
+ }
+ if (!(unwrapped instanceof Slot)) {
+ return null;
+ }
+ Slot slot = (Slot) unwrapped;
+ for (Map.Entry<String, Slot> key : targetKeys.entrySet()) {
+ if (slot.getExprId().equals(key.getValue().getExprId())
+ &&
expression.getDataType().equals(key.getValue().getDataType())) {
+ return key.getKey();
+ }
+ }
+ return null;
+ }
+
+ private AnalysisException invalidNotMatchedKeyCondition() {
+ return new AnalysisException("Connector MERGE with NOT MATCHED
INSERT requires ON to contain "
+ + "only equality predicates for every target primary-key
column and each INSERT "
+ + "to use the corresponding deterministic source
expression");
+ }
+
+ private Alias bindBranchLabel() {
+ Expression targetPresent = new Not(new
org.apache.doris.nereids.trees.expressions.IsNull(
+ findTargetSlot(primaryKeys.get(0))));
+ Expression matched = new NullLiteral(IntegerType.INSTANCE);
+ for (int i = merge.getMatchedClauses().size() - 1; i >= 0; i--) {
+ MergeMatchedClause clause = merge.getMatchedClauses().get(i);
+ if (i != merge.getMatchedClauses().size() - 1 &&
!clause.getCasePredicate().isPresent()) {
+ throw new AnalysisException("Only the last matched clause
may omit its condition");
+ }
+ Expression label = new IntegerLiteral(i);
+ matched = clause.getCasePredicate().isPresent()
+ ? new ShortCircuitIf(clause.getCasePredicate().get(),
label, matched) : label;
+ }
+ Expression notMatched = new NullLiteral(IntegerType.INSTANCE);
+ for (int i = merge.getNotMatchedClauses().size() - 1; i >= 0; i--)
{
+ MergeNotMatchedClause clause =
merge.getNotMatchedClauses().get(i);
+ if (i != merge.getNotMatchedClauses().size() - 1
+ && !clause.getCasePredicate().isPresent()) {
+ throw new AnalysisException("Only the last not matched
clause may omit its condition");
+ }
+ Expression label = new IntegerLiteral(i +
merge.getMatchedClauses().size());
+ notMatched = clause.getCasePredicate().isPresent()
+ ? new ShortCircuitIf(clause.getCasePredicate().get(),
label, notMatched) : label;
+ }
+ return new Alias(analyzer.analyze(
+ new ShortCircuitIf(targetPresent, matched, notMatched)),
BRANCH_LABEL);
+ }
+
+ private List<List<Expression>> buildBranchProjections() {
+ List<List<Expression>> branches = new ArrayList<>();
+ for (MergeMatchedClause clause : merge.getMatchedClauses()) {
+ branches.add(clause.isDelete() ? deleteProjection() :
updateProjection(clause));
+ }
+ for (MergeNotMatchedClause clause : merge.getNotMatchedClauses()) {
+ branches.add(insertProjection(clause));
+ }
+ if (branches.isEmpty()) {
+ throw new AnalysisException("Connector MERGE requires at least
one WHEN clause");
+ }
+ for (List<Expression> branch : branches) {
+ for (int i = 0; i < branch.size(); i++) {
+ branch.set(i, analyzer.analyze(branch.get(i)));
+ }
+ }
+ return branches;
+ }
+
+ private List<Expression> deleteProjection() {
+ List<Expression> output = new ArrayList<>();
+ output.add(new TinyIntLiteral(mode.getDeleteValue()));
+ for (Column column : schema) {
+ output.add(targetSlot(column.getName()));
+ }
+ return output;
+ }
+
+ private List<Expression> updateProjection(MergeMatchedClause clause) {
+ Map<String, Expression> changes =
Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER);
+ for (EqualTo assignment : clause.getAssignments()) {
+ List<String> parts = ((UnboundSlot)
assignment.left()).getNameParts();
+ String name = parts.get(parts.size() - 1);
+ if (changes.put(name, assignment.right()) != null) {
+ throw new AnalysisException("Duplicate column name in
connector MERGE UPDATE: " + name);
+ }
+ }
+ List<Expression> output = new ArrayList<>();
+ output.add(new TinyIntLiteral(mode.getUpdateValue()));
+ for (Column column : schema) {
+ output.add(changes.containsKey(column.getName())
+ ? changes.remove(column.getName()) :
targetSlot(column.getName()));
+ }
+ if (!changes.isEmpty()) {
+ throw new AnalysisException("Unknown column in connector MERGE
UPDATE: "
+ + String.join(", ", changes.keySet()));
+ }
+ return output;
+ }
+
+ private List<Expression> insertProjection(MergeNotMatchedClause
clause) {
+ if (clause.getRow().size() != schema.size()) {
Review Comment:
Addressed in fece43753bd. Changelog MERGE INSERT now expands omitted target
columns from catalog defaults and supports explicit DEFAULT, with
partial-column/default coverage.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ChangelogRowLevelDmlTransform.java:
##########
@@ -0,0 +1,233 @@
+// 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.doris.nereids.trees.plans.commands;
+
+import org.apache.doris.analysis.UserIdentity;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.DatabaseIf;
+import org.apache.doris.catalog.TableIf;
+import org.apache.doris.connector.spi.DorisConnectorException;
+import org.apache.doris.connector.spi.handle.WriteOperation;
+import org.apache.doris.connector.spi.pushdown.ConnectorPredicate;
+import org.apache.doris.connector.spi.write.ConnectorRowChangeStyle;
+import org.apache.doris.connector.spi.write.ConnectorRowLevelDmlRequest;
+import org.apache.doris.datasource.CatalogIf;
+import org.apache.doris.datasource.plugin.PluginDrivenExternalTable;
+import org.apache.doris.mysql.privilege.AccessControllerManager;
+import org.apache.doris.nereids.NereidsPlanner;
+import org.apache.doris.nereids.analyzer.UnboundConnectorTableSink;
+import org.apache.doris.nereids.analyzer.UnboundRelation;
+import org.apache.doris.nereids.analyzer.UnboundSlot;
+import org.apache.doris.nereids.exceptions.AnalysisException;
+import org.apache.doris.nereids.parser.LogicalPlanBuilderAssistant;
+import org.apache.doris.nereids.rules.exploration.join.JoinReorderContext;
+import org.apache.doris.nereids.trees.expressions.EqualTo;
+import org.apache.doris.nereids.trees.expressions.StatementScopeIdGenerator;
+import org.apache.doris.nereids.trees.plans.JoinType;
+import org.apache.doris.nereids.trees.plans.Plan;
+import
org.apache.doris.nereids.trees.plans.commands.info.ConnectorChangelogRowChangeSpec;
+import
org.apache.doris.nereids.trees.plans.commands.insert.BaseExternalTableInsertExecutor;
+import
org.apache.doris.nereids.trees.plans.commands.insert.PluginDrivenInsertExecutor;
+import org.apache.doris.nereids.trees.plans.commands.merge.MergeMatchedClause;
+import org.apache.doris.nereids.trees.plans.logical.LogicalJoin;
+import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
+import org.apache.doris.nereids.trees.plans.logical.LogicalSubQueryAlias;
+import
org.apache.doris.nereids.trees.plans.physical.PhysicalConnectorTableSink;
+import org.apache.doris.nereids.trees.plans.physical.PhysicalSink;
+import org.apache.doris.nereids.util.RelationUtil;
+import org.apache.doris.planner.DataSink;
+import org.apache.doris.planner.PlanFragment;
+import org.apache.doris.qe.ConnectContext;
+
+import com.google.common.collect.ImmutableList;
+
+import java.util.List;
+import java.util.Optional;
+import java.util.Set;
+import java.util.TreeSet;
+
+/** Plans row-level changes as an operation column followed by a complete
table row. */
+public class ChangelogRowLevelDmlTransform implements RowLevelDmlTransform {
+
+ @Override
+ public boolean handles(TableIf table) {
+ if (!(table instanceof PluginDrivenExternalTable)) {
+ return false;
+ }
+ PluginDrivenExternalTable connectorTable = (PluginDrivenExternalTable)
table;
+ if (connectorTable.getConnectorRowChangeStyle() !=
ConnectorRowChangeStyle.CHANGELOG) {
+ return false;
+ }
+ return RowLevelDmlRegistry.supportsAnyRowLevelDml(
+ connectorTable.connectorSupportedWriteOperations());
+ }
+
+ @Override
+ public void checkMode(TableIf table, RowLevelDmlOp op) {
+ WriteOperation operation = op.toWriteOperation();
+ if (!((PluginDrivenExternalTable)
table).connectorSupportedWriteOperations().contains(operation)) {
+ throw new AnalysisException("Connector does not support " +
operation + " operations");
+ }
+ // Statement-specific validation runs in synthesize, where assignments
and MERGE clauses are available.
+ }
+
+ @Override
+ public LogicalPlan synthesize(ConnectContext ctx, RowLevelDmlArgs args,
RowLevelDmlOp op) {
+ PluginDrivenExternalTable table = (PluginDrivenExternalTable)
args.getTable();
+ if (op == RowLevelDmlOp.DELETE && (args.isTempPart() ||
!args.getPartitions().isEmpty())) {
+ throw new AnalysisException(
+ "Connector changelog DELETE does not support partition
name lists; use a WHERE predicate");
+ }
+ validate(ctx, table, args, op);
+ switch (op) {
+ case DELETE:
+ return deletePlan(ctx, args);
+ case UPDATE:
+ return updatePlan(ctx, args);
+ default:
+ return mergePlan(ctx, args);
+ }
+ }
+
+ private LogicalPlan deletePlan(ConnectContext ctx, RowLevelDmlArgs args) {
+ List<String> target = args.getTableAlias() != null
+ ? ImmutableList.of(args.getTableAlias())
+ : RelationUtil.getQualifierName(ctx, args.getNameParts());
+ return new UnboundConnectorTableSink<>(args.getNameParts(),
args.getLogicalQuery(),
+ new ConnectorChangelogRowChangeSpec.Delete(target,
args.shouldDeduplicateTargetRows()));
+ }
+
+ private LogicalPlan updatePlan(ConnectContext ctx, RowLevelDmlArgs args) {
+ for (EqualTo assignment : args.getAssignments()) {
+ UpdateCommand.checkAssignmentColumn(ctx,
+ ((UnboundSlot) assignment.left()).getNameParts(),
+ args.getNameParts(), args.getTableAlias());
+ }
+ List<String> target = args.getTableAlias() != null
+ ? ImmutableList.of(args.getTableAlias())
+ : RelationUtil.getQualifierName(ctx, args.getNameParts());
+ LogicalPlan sink = new
UnboundConnectorTableSink<>(args.getNameParts(), args.getLogicalQuery(),
+ new ConnectorChangelogRowChangeSpec.Update(target,
args.getAssignments()));
+ return args.getCte().isPresent() ? (LogicalPlan)
args.getCte().get().withChildren(sink) : sink;
+ }
+
+ private LogicalPlan mergePlan(ConnectContext ctx, RowLevelDmlArgs args) {
+ for (MergeMatchedClause clause : args.getMatchedClauses()) {
+ for (EqualTo assignment : clause.getAssignments()) {
+ UpdateCommand.checkAssignmentColumn(ctx,
+ ((UnboundSlot) assignment.left()).getNameParts(),
+ args.getTargetNameParts(),
args.getTargetAlias().orElse(null));
+ }
+ }
+ List<String> targetName = args.getTargetAlias().isPresent()
+ ? ImmutableList.of(args.getTargetAlias().get())
+ : RelationUtil.getQualifierName(ctx,
args.getTargetNameParts());
+ ConnectorChangelogRowChangeSpec.Merge spec = new
ConnectorChangelogRowChangeSpec.Merge(
+ targetName, args.getMatchedClauses(),
args.getNotMatchedClauses());
+ LogicalPlan target = LogicalPlanBuilderAssistant.withCheckPolicy(
+ new UnboundRelation(StatementScopeIdGenerator.newRelationId(),
args.getTargetNameParts()));
+ if (args.getTargetAlias().isPresent()) {
+ target = new LogicalSubQueryAlias<>(args.getTargetAlias().get(),
target);
+ }
+ JoinType joinType = args.getNotMatchedClauses().isEmpty()
Review Comment:
Addressed in fece43753bd. The changelog path now reuses
MergeUtils.buildMergeJoin, retaining the established target RIGHT OUTER JOIN
source shape.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalConnectorTableSink.java:
##########
@@ -290,6 +357,12 @@ public PhysicalProperties getRequirePhysicalProperties() {
}
PluginDrivenExternalTable table = (PluginDrivenExternalTable)
targetTable;
+ Optional<ConnectorWriteDistribution> connectorDistribution
+ = table.getConnectorWriteDistribution();
+ if (connectorDistribution.isPresent()) {
+ return toPhysicalProperties(connectorDistribution.get());
Review Comment:
Addressed in fece43753bd. Connector-owned distribution now composes with the
independently required partition-local sort; combined hash-plus-local-sort
behavior is covered.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalConnectorTableSink.java:
##########
@@ -370,4 +446,53 @@ public PhysicalProperties getRequirePhysicalProperties() {
}
return PhysicalProperties.GATHER;
}
+
+ private PhysicalProperties toPhysicalProperties(ConnectorWriteDistribution
distribution) {
+ switch (distribution.getMode()) {
+ case EXECUTION_ANY:
+ return PhysicalProperties.EXECUTION_ANY;
+ case GATHER:
+ return PhysicalProperties.GATHER;
+ case HASH:
+ return PhysicalProperties.createHash(
+ routeExprIds(distribution.getRouteColumns()),
ShuffleType.REQUIRE);
+ case EXTERNAL_UNPARTITIONED:
+ requireExternalWriterRoutingSupport();
+ return PhysicalProperties.EXTERNAL_TABLE_SINK_UNPARTITIONED;
+ case EXTERNAL_HASH:
+ requireExternalWriterRoutingSupport();
+ return new PhysicalProperties(new
DistributionSpecExternalTableSinkHashPartitioned(
Review Comment:
Addressed in fece43753bd. Disabling strict-consistency only relaxes RANDOM
distribution; connector-owned HASH/EXTERNAL_HASH routing remains a correctness
requirement, with strict=false coverage.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalConnectorTableSink.java:
##########
@@ -370,4 +446,53 @@ public PhysicalProperties getRequirePhysicalProperties() {
}
return PhysicalProperties.GATHER;
}
+
+ private PhysicalProperties toPhysicalProperties(ConnectorWriteDistribution
distribution) {
+ switch (distribution.getMode()) {
+ case EXECUTION_ANY:
+ return PhysicalProperties.EXECUTION_ANY;
+ case GATHER:
+ return PhysicalProperties.GATHER;
+ case HASH:
+ return PhysicalProperties.createHash(
+ routeExprIds(distribution.getRouteColumns()),
ShuffleType.REQUIRE);
+ case EXTERNAL_UNPARTITIONED:
+ requireExternalWriterRoutingSupport();
+ return PhysicalProperties.EXTERNAL_TABLE_SINK_UNPARTITIONED;
+ case EXTERNAL_HASH:
+ requireExternalWriterRoutingSupport();
+ return new PhysicalProperties(new
DistributionSpecExternalTableSinkHashPartitioned(
+ routeExprIds(distribution.getRouteColumns()),
+ distribution.getPartitionFunction(),
+ distribution.getPartitionFunctionOptions(),
+ distribution.getWriterAssignment()));
+ default:
+ throw new IllegalStateException("Unsupported connector write
distribution: "
+ + distribution.getMode());
+ }
+ }
+
+ private void requireExternalWriterRoutingSupport() {
+ Preconditions.checkState(Config.be_exec_version
+ >=
DistributionSpecExternalTableSinkHashPartitioned.MIN_BE_EXEC_VERSION,
+ "External table sink distribution requires BE execution
version %s or newer",
+
DistributionSpecExternalTableSinkHashPartitioned.MIN_BE_EXEC_VERSION);
+ }
+
+ private List<ExprId> routeExprIds(List<String> routeColumns) {
+ List<Slot> output = child().getOutput();
+ int offset = hasRowOperationColumn() ? 1 : 0;
+ Preconditions.checkState(boundTargetSchema.size() + offset ==
output.size(),
Review Comment:
Addressed in fece43753bd. Routed slot lookup now uses cols for name-mapped
output and boundTargetSchema only for full-schema positional output, including
the operation-column offset.
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]