From 6ec05b88a365667efe815fff42c8c9c309cecad2 Mon Sep 17 00:00:00 2001 From: wraymo Date: Thu, 16 Oct 2025 21:50:19 -0400 Subject: [PATCH 1/3] add clp_get_json_string --- .../presto/plugin/clp/ClpFunctions.java | 9 +++++ .../clp/optimization/ClpUdfRewriter.java | 23 ++++++++--- .../presto/plugin/clp/TestClpUdfRewriter.java | 39 +++++++++++++++++-- 3 files changed, 62 insertions(+), 9 deletions(-) diff --git a/presto-clp/src/main/java/com/facebook/presto/plugin/clp/ClpFunctions.java b/presto-clp/src/main/java/com/facebook/presto/plugin/clp/ClpFunctions.java index 14495ebb2f41b..9da65562ff26c 100644 --- a/presto-clp/src/main/java/com/facebook/presto/plugin/clp/ClpFunctions.java +++ b/presto-clp/src/main/java/com/facebook/presto/plugin/clp/ClpFunctions.java @@ -19,6 +19,7 @@ import com.facebook.presto.spi.function.ScalarFunction; import com.facebook.presto.spi.function.SqlType; import io.airlift.slice.Slice; +import io.airlift.slice.Slices; public final class ClpFunctions { @@ -97,4 +98,12 @@ public static boolean clpWildcardBoolColumn() { throw new UnsupportedOperationException("CLP_WILDCARD_BOOL_COLUMN is a placeholder function without implementation."); } + + @ScalarFunction(value = "CLP_GET_JSON_STRING", deterministic = false) + @Description("Converts an entire log record into a JSON string.") + @SqlType(StandardTypes.VARCHAR) + public static Slice clpGetJSONString() + { + throw new UnsupportedOperationException("CLP_GET_JSON_STRING is a placeholder function without implementation."); + } } diff --git a/presto-clp/src/main/java/com/facebook/presto/plugin/clp/optimization/ClpUdfRewriter.java b/presto-clp/src/main/java/com/facebook/presto/plugin/clp/optimization/ClpUdfRewriter.java index aff40041e2d58..5237b84845501 100644 --- a/presto-clp/src/main/java/com/facebook/presto/plugin/clp/optimization/ClpUdfRewriter.java +++ b/presto-clp/src/main/java/com/facebook/presto/plugin/clp/optimization/ClpUdfRewriter.java @@ -58,6 +58,7 @@ public final class ClpUdfRewriter implements ConnectorPlanOptimizer { + public static final String JSON_STRING_PLACEHOLDER = "__json_string"; private final FunctionMetadataManager functionManager; public ClpUdfRewriter(FunctionMetadataManager functionManager) @@ -123,7 +124,7 @@ public PlanNode visitProject(ProjectNode node, RewriteContext context) for (Map.Entry entry : node.getAssignments().getMap().entrySet()) { newAssignments.put( entry.getKey(), - rewriteClpUdfs(entry.getValue(), functionManager, variableAllocator)); + rewriteClpUdfs(entry.getValue(), functionManager, variableAllocator, true)); } PlanNode newSource = rewritePlanSubtree(node.getSource()); @@ -148,20 +149,30 @@ public PlanNode visitFilter(FilterNode node, RewriteContext context) * @param expression the input expression to analyze and possibly rewrite * @param functionManager function manager used to resolve function metadata * @param variableAllocator variable allocator used to create new variable references + * @param inProjectNode whether the CLP UDFs are in a {@link ProjectNode} * @return a possibly rewritten {@link RowExpression} with CLP_GET_* calls * replaced */ private RowExpression rewriteClpUdfs( RowExpression expression, FunctionMetadataManager functionManager, - VariableAllocator variableAllocator) + VariableAllocator variableAllocator, + boolean inProjectNode) { // Handle CLP_GET_* function calls if (expression instanceof CallExpression) { CallExpression call = (CallExpression) expression; String functionName = functionManager.getFunctionMetadata(call.getFunctionHandle()).getName().getObjectName().toUpperCase(); - if (functionName.startsWith("CLP_GET_")) { + if (inProjectNode && functionName.equals("CLP_GET_JSON_STRING")) { + VariableReferenceExpression newValue = + new VariableReferenceExpression(expression.getSourceLocation(), JSON_STRING_PLACEHOLDER, expression.getType()); + ClpColumnHandle targetHandle = new ClpColumnHandle(JSON_STRING_PLACEHOLDER, call.getType()); + + globalColumnVarMap.put(targetHandle, newValue); + return newValue; + } + else if (functionName.startsWith("CLP_GET_")) { if (call.getArguments().size() != 1 || !(call.getArguments().get(0) instanceof ConstantExpression)) { throw new PrestoException(CLP_PUSHDOWN_UNSUPPORTED_EXPRESSION, "CLP_GET_* UDF must have a single constant string argument"); @@ -187,7 +198,7 @@ private RowExpression rewriteClpUdfs( // Recurse into arguments List rewrittenArgs = call.getArguments().stream() - .map(arg -> rewriteClpUdfs(arg, functionManager, variableAllocator)) + .map(arg -> rewriteClpUdfs(arg, functionManager, variableAllocator, inProjectNode)) .collect(toImmutableList()); return new CallExpression(call.getDisplayName(), call.getFunctionHandle(), call.getType(), rewrittenArgs); @@ -198,7 +209,7 @@ private RowExpression rewriteClpUdfs( SpecialFormExpression special = (SpecialFormExpression) expression; List rewrittenArgs = special.getArguments().stream() - .map(arg -> rewriteClpUdfs(arg, functionManager, variableAllocator)) + .map(arg -> rewriteClpUdfs(arg, functionManager, variableAllocator, inProjectNode)) .collect(toImmutableList()); return new SpecialFormExpression(special.getSourceLocation(), special.getForm(), special.getType(), rewrittenArgs); @@ -298,7 +309,7 @@ private TableScanNode buildNewTableScanNode(TableScanNode node) */ private FilterNode buildNewFilterNode(FilterNode node) { - RowExpression newPredicate = rewriteClpUdfs(node.getPredicate(), functionManager, variableAllocator); + RowExpression newPredicate = rewriteClpUdfs(node.getPredicate(), functionManager, variableAllocator, false); PlanNode newSource = rewritePlanSubtree(node.getSource()); return new FilterNode(node.getSourceLocation(), idAllocator.getNextId(), newSource, newPredicate); } diff --git a/presto-clp/src/test/java/com/facebook/presto/plugin/clp/TestClpUdfRewriter.java b/presto-clp/src/test/java/com/facebook/presto/plugin/clp/TestClpUdfRewriter.java index b82866f3d8dd9..a6b6bc118ffee 100644 --- a/presto-clp/src/test/java/com/facebook/presto/plugin/clp/TestClpUdfRewriter.java +++ b/presto-clp/src/test/java/com/facebook/presto/plugin/clp/TestClpUdfRewriter.java @@ -72,6 +72,7 @@ import static com.facebook.presto.plugin.clp.metadata.ClpSchemaTreeNodeType.Float; import static com.facebook.presto.plugin.clp.metadata.ClpSchemaTreeNodeType.Integer; import static com.facebook.presto.plugin.clp.metadata.ClpSchemaTreeNodeType.VarString; +import static com.facebook.presto.plugin.clp.optimization.ClpUdfRewriter.JSON_STRING_PLACEHOLDER; import static com.facebook.presto.sql.planner.assertions.MatchResult.NO_MATCH; import static com.facebook.presto.sql.planner.assertions.MatchResult.match; import static com.facebook.presto.sql.planner.assertions.PlanMatchPattern.anyTree; @@ -140,7 +141,7 @@ public void tearDown() } @Test - public void testScanFilter() + public void testClpGetScanFilter() { TransactionId transactionId = localQueryRunner.getTransactionManager().beginTransaction(false); Session session = testSessionBuilder().setCatalog("clp").setSchema("default").setTransactionId(transactionId).build(); @@ -179,7 +180,7 @@ public void testScanFilter() } @Test - public void testScanProject() + public void testClpGetScanProject() { TransactionId transactionId = localQueryRunner.getTransactionManager().beginTransaction(false); Session session = testSessionBuilder().setCatalog("clp").setSchema("default").setTransactionId(transactionId).build(); @@ -229,7 +230,7 @@ public void testScanProject() } @Test - public void testScanProjectFilter() + public void testClpGetScanProjectFilter() { TransactionId transactionId = localQueryRunner.getTransactionManager().beginTransaction(false); Session session = testSessionBuilder().setCatalog("clp").setSchema("default").setTransactionId(transactionId).build(); @@ -265,6 +266,38 @@ public void testScanProjectFilter() city)))))); } + @Test + public void testClpGetJsonString() + { + TransactionId transactionId = localQueryRunner.getTransactionManager().beginTransaction(false); + Session session = testSessionBuilder().setCatalog("clp").setSchema("default").setTransactionId(transactionId).build(); + + Plan plan = localQueryRunner.createPlan( + session, + "SELECT CLP_GET_JSON_STRING() from test WHERE CLP_GET_BIGINT('user_id') = 0", + WarningCollector.NOOP); + ClpUdfRewriter udfRewriter = new ClpUdfRewriter(functionAndTypeManager); + PlanNode optimizedPlan = udfRewriter.optimize(plan.getRoot(), session.toConnectorSession(), variableAllocator, planNodeIdAllocator); + ClpComputePushDown optimizer = new ClpComputePushDown(functionAndTypeManager, functionResolution, splitFilterProvider); + optimizedPlan = optimizer.optimize(optimizedPlan, session.toConnectorSession(), variableAllocator, planNodeIdAllocator); + + PlanAssert.assertPlan( + session, + localQueryRunner.getMetadata(), + (node, sourceStats, lookup, s, types) -> PlanNodeStatsEstimate.unknown(), + new Plan(optimizedPlan, plan.getTypes(), StatsAndCosts.empty()), + anyTree( + project( + ImmutableMap.of( + "clp_get_json_string", + PlanMatchPattern.expression(JSON_STRING_PLACEHOLDER)), + ClpTableScanMatcher.clpTableScanPattern( + new ClpTableLayoutHandle(table, Optional.of("user_id: 0"), Optional.empty()), + ImmutableSet.of( + new ClpColumnHandle("user_id", BIGINT), + new ClpColumnHandle(JSON_STRING_PLACEHOLDER, VARCHAR)))))); + } + private static final class ClpTableScanMatcher implements Matcher { From 497d6e4659837e9c68c25fa51837a07e1c1d6243 Mon Sep 17 00:00:00 2001 From: wraymo Date: Sat, 18 Oct 2025 14:20:58 -0400 Subject: [PATCH 2/3] fix --- .../java/com/facebook/presto/plugin/clp/ClpFunctions.java | 1 - .../presto/plugin/clp/optimization/ClpUdfRewriter.java | 6 ++++-- 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/presto-clp/src/main/java/com/facebook/presto/plugin/clp/ClpFunctions.java b/presto-clp/src/main/java/com/facebook/presto/plugin/clp/ClpFunctions.java index 9da65562ff26c..7ced54299e7b7 100644 --- a/presto-clp/src/main/java/com/facebook/presto/plugin/clp/ClpFunctions.java +++ b/presto-clp/src/main/java/com/facebook/presto/plugin/clp/ClpFunctions.java @@ -19,7 +19,6 @@ import com.facebook.presto.spi.function.ScalarFunction; import com.facebook.presto.spi.function.SqlType; import io.airlift.slice.Slice; -import io.airlift.slice.Slices; public final class ClpFunctions { diff --git a/presto-clp/src/main/java/com/facebook/presto/plugin/clp/optimization/ClpUdfRewriter.java b/presto-clp/src/main/java/com/facebook/presto/plugin/clp/optimization/ClpUdfRewriter.java index 5237b84845501..75d0a66cc02c7 100644 --- a/presto-clp/src/main/java/com/facebook/presto/plugin/clp/optimization/ClpUdfRewriter.java +++ b/presto-clp/src/main/java/com/facebook/presto/plugin/clp/optimization/ClpUdfRewriter.java @@ -165,8 +165,10 @@ private RowExpression rewriteClpUdfs( String functionName = functionManager.getFunctionMetadata(call.getFunctionHandle()).getName().getObjectName().toUpperCase(); if (inProjectNode && functionName.equals("CLP_GET_JSON_STRING")) { - VariableReferenceExpression newValue = - new VariableReferenceExpression(expression.getSourceLocation(), JSON_STRING_PLACEHOLDER, expression.getType()); + VariableReferenceExpression newValue = variableAllocator.newVariable( + expression.getSourceLocation(), + JSON_STRING_PLACEHOLDER, + call.getType()); ClpColumnHandle targetHandle = new ClpColumnHandle(JSON_STRING_PLACEHOLDER, call.getType()); globalColumnVarMap.put(targetHandle, newValue); From 144a4493c8a4aa40d906dee18aa87311b43ebb49 Mon Sep 17 00:00:00 2001 From: wraymo Date: Mon, 27 Oct 2025 10:04:52 -0400 Subject: [PATCH 3/3] advance velox --- presto-native-execution/velox | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/presto-native-execution/velox b/presto-native-execution/velox index cf895960f83ec..425b66f6cb39b 160000 --- a/presto-native-execution/velox +++ b/presto-native-execution/velox @@ -1 +1 @@ -Subproject commit cf895960f83ece5621ba3dd40f68f322a4bab570 +Subproject commit 425b66f6cb39b9e20cf6801af22aa6126562c999