diff --git a/samza-sql/src/main/java/org/apache/samza/sql/interfaces/UdfMetadata.java b/samza-sql/src/main/java/org/apache/samza/sql/interfaces/UdfMetadata.java index 9288ce7e63..c06f60e5d1 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/interfaces/UdfMetadata.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/interfaces/UdfMetadata.java @@ -43,7 +43,8 @@ public class UdfMetadata { public UdfMetadata(String name, String description, Method udfMethod, Config udfConfig, List arguments, SamzaSqlFieldType returnType, boolean disableArgCheck) { - this.name = name; + // Udfs are case insensitive + this.name = name.toUpperCase(); this.description = description; this.udfMethod = udfMethod; this.udfConfig = udfConfig; diff --git a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlUdfOperatorTable.java b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlUdfOperatorTable.java index 6ee10f85e5..151bb02914 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlUdfOperatorTable.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlUdfOperatorTable.java @@ -62,7 +62,12 @@ private SqlOperator getSqlOperator(SamzaSqlScalarFunctionImpl scalarFunction) { @Override public void lookupOperatorOverloads(SqlIdentifier opName, SqlFunctionCategory category, SqlSyntax syntax, List operatorList) { - operatorTable.lookupOperatorOverloads(opName, category, syntax, operatorList); + SqlIdentifier upperCaseOpName = opName; + // Only udfs are case insensitive + if (category != null && category.equals(SqlFunctionCategory.USER_DEFINED_FUNCTION)) { + upperCaseOpName = new SqlIdentifier(opName.names.get(0).toUpperCase(), opName.getComponentParserPosition(0)); + } + operatorTable.lookupOperatorOverloads(upperCaseOpName, category, syntax, operatorList); } @Override diff --git a/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java b/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java index aef79265b3..59f10ade2a 100644 --- a/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java +++ b/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java @@ -504,7 +504,7 @@ public void testEndToEndUdf() throws Exception { TestAvroSystemFactory.messages.clear(); Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); String sql1 = "Insert into testavro.outputTopic(id, long_value) " - + "select id, MyTest(id) as long_value from testavro.SIMPLE1"; + + "select id, MYTest(id) as long_value from testavro.SIMPLE1"; List sqlStmts = Collections.singletonList(sql1); staticConfigs.put(SamzaSqlApplicationConfig.CFG_SQL_STMTS_JSON, JsonUtil.toJson(sqlStmts)); runApplication(new MapConfig(staticConfigs));