[ 
https://issues.apache.org/jira/browse/DRILL-8552?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18105965#comment-18105965
 ] 

ASF GitHub Bot commented on DRILL-8552:
---------------------------------------

cgivre commented on code in PR #3067:
URL: https://github.com/apache/drill/pull/3067#discussion_r3814420299


##########
contrib/storage-accumulo/src/main/java/org/apache/drill/exec/store/accumulo/AccumuloFilterBuilder.java:
##########
@@ -0,0 +1,326 @@
+/*
+ * 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.drill.exec.store.accumulo;
+
+import java.util.Arrays;
+import java.util.List;
+
+import org.apache.drill.common.FunctionNames;
+import org.apache.drill.common.expression.BooleanOperator;
+import org.apache.drill.common.expression.FunctionCall;
+import org.apache.drill.common.expression.LogicalExpression;
+import org.apache.drill.common.expression.SchemaPath;
+import org.apache.drill.common.expression.visitors.AbstractExprVisitor;
+
+/**
+ * Builds Accumulo scan specifications from Drill filter expressions.
+ *
+ * <p>This class converts Drill's LogicalExpression filter representation into
+ * Accumulo scan parameters (start row, stop row). It focuses on row key
+ * predicates since those can be efficiently pushed down to Accumulo's scan 
range.</p>
+ *
+ * <p>Supported predicates on row_key:</p>
+ * <ul>
+ *   <li>row_key = 'value' → exact range</li>
+ *   <li>row_key > 'value' → start row (exclusive)</li>
+ *   <li>row_key >= 'value' → start row (inclusive)</li>
+ *   <li>row_key < 'value' → stop row (exclusive)</li>
+ *   <li>row_key <= 'value' → stop row (inclusive)</li>
+ *   <li>AND combinations → intersect ranges</li>
+ *   <li>OR combinations → union ranges (if contiguous)</li>
+ * </ul>
+ */
+public class AccumuloFilterBuilder
+    extends AbstractExprVisitor<AccumuloScanSpec, Void, RuntimeException>
+    implements DrillAccumuloConstants {
+
+  private final AccumuloGroupScan groupScan;
+  private final LogicalExpression filterExpression;
+  private boolean allExpressionsConverted = true;
+
+  public AccumuloFilterBuilder(AccumuloGroupScan groupScan, LogicalExpression 
filterExpression) {
+    this.groupScan = groupScan;
+    this.filterExpression = filterExpression;
+  }
+
+  /**
+   * Parses the filter expression and returns an updated scan specification.
+   *
+   * @return the scan spec with row key ranges, or null if no filters can be 
pushed
+   */
+  public AccumuloScanSpec parseTree() {
+    AccumuloScanSpec parsedSpec = filterExpression.accept(this, null);
+    if (parsedSpec != null) {
+      // Merge with existing scan spec
+      parsedSpec = mergeScanSpecs(FunctionNames.AND, groupScan.getScanSpec(), 
parsedSpec);
+    }
+    return parsedSpec;
+  }
+
+  /**
+   * Returns true if all filter expressions were converted to Accumulo scan 
parameters.
+   * If false, the filter operator should remain in the plan for client-side 
filtering.
+   */
+  public boolean isAllExpressionsConverted() {
+    return allExpressionsConverted;
+  }
+
+  @Override
+  public AccumuloScanSpec visitUnknown(LogicalExpression e, Void value) throws 
RuntimeException {
+    allExpressionsConverted = false;
+    return null;
+  }
+
+  @Override
+  public AccumuloScanSpec visitBooleanOperator(BooleanOperator op, Void value)
+      throws RuntimeException {
+    return visitFunctionCall(op, value);
+  }
+
+  @Override
+  public AccumuloScanSpec visitFunctionCall(FunctionCall call, Void value)
+      throws RuntimeException {
+    AccumuloScanSpec nodeScanSpec = null;
+    String functionName = call.getName();
+    List<LogicalExpression> args = call.args();
+
+    if (AccumuloCompareFunctionsProcessor.isCompareFunction(functionName)) {
+      AccumuloCompareFunctionsProcessor processor =
+          
AccumuloCompareFunctionsProcessor.createFunctionsProcessorInstance(call);
+      if (processor.isSuccess()) {
+        nodeScanSpec = createScanSpecFromComparison(processor);
+      }
+    } else {
+      switch (functionName) {
+        case FunctionNames.AND:
+        case FunctionNames.OR:
+          AccumuloScanSpec firstScanSpec = args.get(0).accept(this, null);
+          for (int i = 1; i < args.size(); ++i) {
+            AccumuloScanSpec nextScanSpec = args.get(i).accept(this, null);
+            if (firstScanSpec != null && nextScanSpec != null) {
+              nodeScanSpec = mergeScanSpecs(functionName, firstScanSpec, 
nextScanSpec);
+            } else {
+              allExpressionsConverted = false;
+              if (FunctionNames.AND.equals(functionName)) {
+                // For AND, keep whichever spec we have
+                nodeScanSpec = firstScanSpec == null ? nextScanSpec : 
firstScanSpec;
+              }
+              // For OR, if either is null we can't push down the whole OR
+            }
+            firstScanSpec = nodeScanSpec;
+          }
+          break;
+        default:
+          // Unknown function
+          break;
+      }
+    }
+
+    if (nodeScanSpec == null) {
+      allExpressionsConverted = false;
+    }
+
+    return nodeScanSpec;
+  }
+
+  /**
+   * Creates a scan spec from a comparison processor result.
+   */
+  private AccumuloScanSpec createScanSpecFromComparison(
+      AccumuloCompareFunctionsProcessor processor) {
+
+    String functionName = processor.getFunctionName();
+    SchemaPath field = processor.getPath();
+    byte[] fieldValue = processor.getValue();
+
+    // Only handle row_key predicates for now
+    boolean isRowKey = field.getRootSegmentPath().equalsIgnoreCase(ROW_KEY);
+    if (!isRowKey) {
+      // Column predicates require iterators - not supported in Option A
+      return null;
+    }
+
+    byte[] startRow = null;
+    byte[] stopRow = null;
+    boolean startRowInclusive = true;
+    boolean stopRowInclusive = false;
+
+    switch (functionName) {
+      case FunctionNames.EQ:
+        // row_key = 'value' → scan exactly that row
+        startRow = fieldValue;
+        // Stop row should be just after the value
+        stopRow = Arrays.copyOf(fieldValue, fieldValue.length + 1);
+        startRowInclusive = true;
+        stopRowInclusive = false;
+        break;
+
+      case FunctionNames.NE:
+        // row_key != 'value' → can't efficiently push down (would need full 
scan minus one row)
+        return null;
+
+      case FunctionNames.GE:
+        // row_key >= 'value' → start at value (inclusive)
+        startRow = fieldValue;
+        startRowInclusive = true;
+        break;
+
+      case FunctionNames.GT:
+        // row_key > 'value' → start just after value
+        startRow = Arrays.copyOf(fieldValue, fieldValue.length + 1);
+        startRowInclusive = true;
+        break;
+
+      case FunctionNames.LE:
+        // row_key <= 'value' → stop just after value
+        stopRow = Arrays.copyOf(fieldValue, fieldValue.length + 1);
+        stopRowInclusive = false;
+        break;
+
+      case FunctionNames.LT:
+        // row_key < 'value' → stop at value (exclusive)
+        stopRow = fieldValue;
+        stopRowInclusive = false;
+        break;
+
+      default:
+        return null;
+    }
+
+    return new AccumuloScanSpec(
+        groupScan.getTableName(),
+        startRow,
+        stopRow,
+        startRowInclusive,
+        stopRowInclusive,
+        groupScan.getScanSpec().getColumns(),
+        null, // No filter expression needed when using row ranges
+        groupScan.getScanSpec().getLimit(),
+        groupScan.getScanSpec().isUseSortedScanner(),
+        groupScan.getScanSpec().isSortDescending());
+  }
+
+  /**
+   * Merges two scan specs using AND or OR logic.
+   */
+  private AccumuloScanSpec mergeScanSpecs(
+      String functionName,
+      AccumuloScanSpec leftSpec,
+      AccumuloScanSpec rightSpec) {
+
+    byte[] startRow = null;
+    byte[] stopRow = null;
+    boolean startRowInclusive = true;
+    boolean stopRowInclusive = false;
+
+    switch (functionName) {
+      case FunctionNames.AND:
+        // AND: Take the intersection (max of starts, min of stops)
+        startRow = maxOfStartRows(leftSpec.getStartRow(), 
rightSpec.getStartRow());
+        stopRow = minOfStopRows(leftSpec.getStopRow(), rightSpec.getStopRow());
+        break;
+
+      case FunctionNames.OR:
+        // OR: Take the union (min of starts, max of stops)

Review Comment:
   Good catch, this was a real correctness bug. The union of the two ranges is 
a superset of the disjunction, but the builder still reported 
`allExpressionsConverted == true`, so `AccumuloPushFilterIntoScan` dropped the 
filter and the scan returned everything between 'a' and 'z'.
   
   Fixed in a7b5c66: an OR merge now sets `allExpressionsConverted = false`, so 
the range still narrows the scan but the Drill filter stays in the plan to 
discard the rows in between.
   
   While in there I also fixed a second bug in the same path: the OR union used 
the AND helpers, which treat a null bound as "take the other side" rather than 
"unbounded". That made `row_key < 'x' OR row_key > 'y'` collapse into an 
inverted range. Union now propagates null (unbounded).
   
   Added two integration tests for this — `testFilterOnRowKeyDisjunction` and 
`testFilterOnRowKeyDisjunctionOfRanges`. Both fail on the old code (the first 
returned 9 rows instead of 2) and pass now.



##########
contrib/storage-accumulo/src/main/resources/drill-module.conf:
##########
@@ -0,0 +1,36 @@
+#
+# 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.
+#
+
+# This file tells Drill to consider this module when class path scanning.
+# This file can also include any supplementary configuration information.
+# This file is in HOCON format, see 
https://github.com/typesafehub/config/blob/master/HOCON.md for more information.
+
+drill: {
+  classpath.scanning: {
+    packages += "org.apache.drill.exec.store.accumulo"
+  }
+
+  exec: {
+    accumulo.scan: {
+      # Number of rows to sample for schema inference
+      samplerows.count: 100,

Review Comment:
   Right, those were leftovers — the values are hardcoded in Java 
(`DrillAccumuloConstants.DEFAULT_BATCH_SIZE` and 
`DrillAccumuloTable.COLUMN_FAMILY_SAMPLE_SIZE`) and nothing ever read the 
config keys. Removed the `exec.accumulo.scan` block in a7b5c66, and the record 
reader now uses `DEFAULT_BATCH_SIZE` instead of its own duplicate 4000.





> Add Storage Plugin for Apache Accumulo
> --------------------------------------
>
>                 Key: DRILL-8552
>                 URL: https://issues.apache.org/jira/browse/DRILL-8552
>             Project: Apache Drill
>          Issue Type: New Feature
>          Components: Storage - Accumulo
>    Affects Versions: 1.22.0
>            Reporter: Charles Givre
>            Assignee: Charles Givre
>            Priority: Major
>             Fix For: 1.23.0
>
>




--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to