[
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)