Skip to content

Commit

Permalink
[fix](join) Disable scan sharing for NAAJ (apache#39480)
Browse files Browse the repository at this point in the history
  • Loading branch information
Gabriel39 authored Aug 19, 2024
1 parent 0b5ee48 commit 2ede4d7
Show file tree
Hide file tree
Showing 2 changed files with 21 additions and 19 deletions.
20 changes: 11 additions & 9 deletions fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
Original file line number Diff line number Diff line change
Expand Up @@ -1876,10 +1876,11 @@ protected void computeFragmentHosts() throws Exception {
if ((isColocateFragment(fragment, fragment.getPlanRoot())
&& fragmentIdToSeqToAddressMap.containsKey(fragment.getFragmentId())
&& fragmentIdToSeqToAddressMap.get(fragment.getFragmentId()).size() > 0)) {
computeColocateJoinInstanceParam(fragment.getFragmentId(), parallelExecInstanceNum, params);
computeColocateJoinInstanceParam(fragment.getFragmentId(), parallelExecInstanceNum, params,
fragment.hasNullAwareLeftAntiJoin());
} else if (bucketShuffleJoinController.isBucketShuffleJoin(fragment.getFragmentId().asInt())) {
bucketShuffleJoinController.computeInstanceParam(fragment.getFragmentId(),
parallelExecInstanceNum, params);
parallelExecInstanceNum, params, fragment.hasNullAwareLeftAntiJoin());
} else {
// case A
for (Entry<TNetworkAddress, Map<Integer, List<TScanRangeParams>>> entry : fragmentExecParamsMap.get(
Expand All @@ -1904,7 +1905,8 @@ protected void computeFragmentHosts() throws Exception {
int expectedInstanceNum = Math.min(parallelExecInstanceNum,
leftMostNode.getNumInstances());
boolean forceToLocalShuffle = context != null
&& context.getSessionVariable().isForceToLocalShuffle();
&& context.getSessionVariable().isForceToLocalShuffle()
&& !fragment.hasNullAwareLeftAntiJoin();
boolean ignoreStorageDataDistribution = forceToLocalShuffle || (node.isPresent()
&& node.get().ignoreStorageDataDistribution(context, addressToBackendID.size())
&& useNereids);
Expand Down Expand Up @@ -2072,9 +2074,9 @@ private List<TScanRangeParams> findOrInsert(Map<Integer, List<TScanRangeParams>>
}

private void computeColocateJoinInstanceParam(PlanFragmentId fragmentId,
int parallelExecInstanceNum, FragmentExecParams params) {
int parallelExecInstanceNum, FragmentExecParams params, boolean hasNullAwareLeftAntiJoin) {
assignScanRanges(fragmentId, parallelExecInstanceNum, params, fragmentIdTobucketSeqToScanRangeMap,
fragmentIdToSeqToAddressMap, fragmentIdToScanNodeIds);
fragmentIdToSeqToAddressMap, fragmentIdToScanNodeIds, hasNullAwareLeftAntiJoin);
}

private Map<TNetworkAddress, Long> getReplicaNumPerHostForOlapTable() {
Expand Down Expand Up @@ -2689,16 +2691,16 @@ private void computeScanRangeAssignmentByBucket(
}

private void computeInstanceParam(PlanFragmentId fragmentId,
int parallelExecInstanceNum, FragmentExecParams params) {
int parallelExecInstanceNum, FragmentExecParams params, boolean hasNullAwareLeftAntiJoin) {
assignScanRanges(fragmentId, parallelExecInstanceNum, params, fragmentIdBucketSeqToScanRangeMap,
fragmentIdToSeqToAddressMap, fragmentIdToScanNodeIds);
fragmentIdToSeqToAddressMap, fragmentIdToScanNodeIds, hasNullAwareLeftAntiJoin);
}
}

private void assignScanRanges(PlanFragmentId fragmentId, int parallelExecInstanceNum, FragmentExecParams params,
Map<PlanFragmentId, BucketSeqToScanRange> fragmentIdBucketSeqToScanRangeMap,
Map<PlanFragmentId, Map<Integer, TNetworkAddress>> curFragmentIdToSeqToAddressMap,
Map<PlanFragmentId, Set<Integer>> fragmentIdToScanNodeIds) {
Map<PlanFragmentId, Set<Integer>> fragmentIdToScanNodeIds, boolean hasNullAwareLeftAntiJoin) {
Map<Integer, TNetworkAddress> bucketSeqToAddress = curFragmentIdToSeqToAddressMap.get(fragmentId);
BucketSeqToScanRange bucketSeqToScanRange = fragmentIdBucketSeqToScanRangeMap.get(fragmentId);
Set<Integer> scanNodeIds = fragmentIdToScanNodeIds.get(fragmentId);
Expand Down Expand Up @@ -2732,7 +2734,7 @@ private void assignScanRanges(PlanFragmentId fragmentId, int parallelExecInstanc
* 2. Use Nereids planner.
*/
boolean forceToLocalShuffle = context != null
&& context.getSessionVariable().isForceToLocalShuffle();
&& context.getSessionVariable().isForceToLocalShuffle() && !hasNullAwareLeftAntiJoin;
boolean ignoreStorageDataDistribution = forceToLocalShuffle || (scanNodes.stream()
.allMatch(node -> node.ignoreStorageDataDistribution(context,
addressToBackendID.size()))
Expand Down
20 changes: 10 additions & 10 deletions fe/fe-core/src/test/java/org/apache/doris/qe/CoordinatorTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -123,7 +123,7 @@ public void testComputeColocateJoinInstanceParam() {
Deencapsulation.setField(coordinator, "fragmentIdTobucketSeqToScanRangeMap", fragmentIdBucketSeqToScanRangeMap);

FragmentExecParams params = new FragmentExecParams(null);
Deencapsulation.invoke(coordinator, "computeColocateJoinInstanceParam", planFragmentId, 1, params);
Deencapsulation.invoke(coordinator, "computeColocateJoinInstanceParam", planFragmentId, 1, params, false);
Assert.assertEquals(1, params.instanceExecParams.size());

// check whether one instance have 3 tablet to scan
Expand All @@ -134,15 +134,15 @@ public void testComputeColocateJoinInstanceParam() {
}

params = new FragmentExecParams(null);
Deencapsulation.invoke(coordinator, "computeColocateJoinInstanceParam", planFragmentId, 2, params);
Deencapsulation.invoke(coordinator, "computeColocateJoinInstanceParam", planFragmentId, 2, params, false);
Assert.assertEquals(2, params.instanceExecParams.size());

params = new FragmentExecParams(null);
Deencapsulation.invoke(coordinator, "computeColocateJoinInstanceParam", planFragmentId, 3, params);
Deencapsulation.invoke(coordinator, "computeColocateJoinInstanceParam", planFragmentId, 3, params, false);
Assert.assertEquals(3, params.instanceExecParams.size());

params = new FragmentExecParams(null);
Deencapsulation.invoke(coordinator, "computeColocateJoinInstanceParam", planFragmentId, 5, params);
Deencapsulation.invoke(coordinator, "computeColocateJoinInstanceParam", planFragmentId, 5, params, false);
Assert.assertEquals(3, params.instanceExecParams.size());
}

Expand Down Expand Up @@ -324,7 +324,7 @@ public void testColocateJoinAssignment() {
PlanFragment fragment = new PlanFragment(planFragmentId, olapScanNode,
new DataPartition(TPartitionType.UNPARTITIONED));
FragmentExecParams params = new FragmentExecParams(fragment);
Deencapsulation.invoke(coordinator, "computeColocateJoinInstanceParam", planFragmentId, 1, params);
Deencapsulation.invoke(coordinator, "computeColocateJoinInstanceParam", planFragmentId, 1, params, false);
StringBuilder sb = new StringBuilder();
params.appendTo(sb);
Assert.assertTrue(sb.toString().contains("range=[id1,range=[]]"));
Expand Down Expand Up @@ -452,19 +452,19 @@ public void testComputeBucketShuffleJoinInstanceParam() {
Deencapsulation.setField(bucketShuffleJoinController, "fragmentIdBucketSeqToScanRangeMap", fragmentIdBucketSeqToScanRangeMap);

FragmentExecParams params = new FragmentExecParams(null);
Deencapsulation.invoke(bucketShuffleJoinController, "computeInstanceParam", planFragmentId, 1, params);
Deencapsulation.invoke(bucketShuffleJoinController, "computeInstanceParam", planFragmentId, 1, params, false);
Assert.assertEquals(1, params.instanceExecParams.size());

params = new FragmentExecParams(null);
Deencapsulation.invoke(bucketShuffleJoinController, "computeInstanceParam", planFragmentId, 2, params);
Deencapsulation.invoke(bucketShuffleJoinController, "computeInstanceParam", planFragmentId, 2, params, false);
Assert.assertEquals(2, params.instanceExecParams.size());

params = new FragmentExecParams(null);
Deencapsulation.invoke(bucketShuffleJoinController, "computeInstanceParam", planFragmentId, 3, params);
Deencapsulation.invoke(bucketShuffleJoinController, "computeInstanceParam", planFragmentId, 3, params, false);
Assert.assertEquals(3, params.instanceExecParams.size());

params = new FragmentExecParams(null);
Deencapsulation.invoke(bucketShuffleJoinController, "computeInstanceParam", planFragmentId, 5, params);
Deencapsulation.invoke(bucketShuffleJoinController, "computeInstanceParam", planFragmentId, 5, params, false);
Assert.assertEquals(3, params.instanceExecParams.size());
}

Expand Down Expand Up @@ -506,7 +506,7 @@ public void testBucketShuffleAssignment() {
new DataPartition(TPartitionType.UNPARTITIONED));

FragmentExecParams params = new FragmentExecParams(fragment);
Deencapsulation.invoke(bucketShuffleJoinController, "computeInstanceParam", planFragmentId, 1, params);
Deencapsulation.invoke(bucketShuffleJoinController, "computeInstanceParam", planFragmentId, 1, params, false);
Assert.assertEquals(1, params.instanceExecParams.size());
StringBuilder sb = new StringBuilder();
params.appendTo(sb);
Expand Down

0 comments on commit 2ede4d7

Please sign in to comment.