From 2e482cdb2463070374d130f718df2d300f640e8d Mon Sep 17 00:00:00 2001 From: Cameron Lee Date: Tue, 20 Aug 2019 14:41:09 -0700 Subject: [PATCH 1/2] SAMZA-2306: Use in-memory system for SQL tests in samza-test --- .../InMemoryIntegrationTestHarness.java | 63 +++++++++++++++++++ .../SamzaSqlIntegrationTestHarness.java | 7 ++- .../test/samzasql/TestSamzaSqlEndToEnd.java | 13 +--- 3 files changed, 69 insertions(+), 14 deletions(-) create mode 100644 samza-test/src/test/java/org/apache/samza/test/harness/InMemoryIntegrationTestHarness.java diff --git a/samza-test/src/test/java/org/apache/samza/test/harness/InMemoryIntegrationTestHarness.java b/samza-test/src/test/java/org/apache/samza/test/harness/InMemoryIntegrationTestHarness.java new file mode 100644 index 0000000000..f120ac3f35 --- /dev/null +++ b/samza-test/src/test/java/org/apache/samza/test/harness/InMemoryIntegrationTestHarness.java @@ -0,0 +1,63 @@ +/* + * 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.samza.test.harness; + +import java.util.HashMap; +import java.util.Map; +import java.util.Optional; +import org.apache.commons.lang.RandomStringUtils; +import org.apache.samza.config.Config; +import org.apache.samza.config.InMemorySystemConfig; +import org.apache.samza.config.JobConfig; +import org.apache.samza.config.MapConfig; +import org.apache.samza.config.SystemConfig; +import org.apache.samza.context.ExternalContext; +import org.apache.samza.runtime.ApplicationRunner; +import org.apache.samza.system.inmemory.InMemorySystemFactory; + + +/** + * Provides helpers for configuring an in-memory system to be used for tests and executing those tests. + * + * This is somewhat based on {@link IntegrationTestHarness}, but it avoids using Kafka/Zookeeper. + */ +public class InMemoryIntegrationTestHarness { + protected static final String IN_MEMORY = "inmemory"; + + protected Config baseInMemorySystemConfigs() { + Map configMap = new HashMap<>(); + configMap.put(String.format(SystemConfig.SYSTEM_FACTORY_FORMAT, IN_MEMORY), InMemorySystemFactory.class.getName()); + configMap.put(InMemorySystemConfig.INMEMORY_SCOPE, RandomStringUtils.random(10, true, true)); + configMap.put(JobConfig.JOB_DEFAULT_SYSTEM, IN_MEMORY); + return new MapConfig(configMap); + } + + protected void executeRun(ApplicationRunner applicationRunner, Config config) { + applicationRunner.run(buildExternalContext(config).orElse(null)); + } + + private Optional buildExternalContext(Config config) { + /* + * By default, use an empty ExternalContext here. In a custom fork of Samza, this can be implemented to pass + * a non-empty ExternalContext. Only config should be used to build the external context. In the future, components + * like the application descriptor may not be available. + */ + return Optional.empty(); + } +} diff --git a/samza-test/src/test/java/org/apache/samza/test/samzasql/SamzaSqlIntegrationTestHarness.java b/samza-test/src/test/java/org/apache/samza/test/samzasql/SamzaSqlIntegrationTestHarness.java index d41418a516..d32d82e6c1 100644 --- a/samza-test/src/test/java/org/apache/samza/test/samzasql/SamzaSqlIntegrationTestHarness.java +++ b/samza-test/src/test/java/org/apache/samza/test/samzasql/SamzaSqlIntegrationTestHarness.java @@ -29,11 +29,11 @@ import org.apache.samza.sql.util.SamzaSqlTestConfig; import org.apache.samza.system.MockSystemFactory; import org.apache.samza.system.SystemStreamPartition; -import org.apache.samza.test.harness.IntegrationTestHarness; +import org.apache.samza.test.harness.InMemoryIntegrationTestHarness; import org.apache.samza.util.CoordinatorStreamUtil; -public class SamzaSqlIntegrationTestHarness extends IntegrationTestHarness { +public class SamzaSqlIntegrationTestHarness extends InMemoryIntegrationTestHarness { public static final String MOCK_METADATA_SYSTEM = "mockmetadatasystem"; @@ -45,6 +45,9 @@ protected void runApplication(Config config) { HashMap mapConfig = new HashMap<>(); mapConfig.put(JobConfig.JOB_COORDINATOR_SYSTEM, MOCK_METADATA_SYSTEM); mapConfig.put(String.format(SystemConfig.SYSTEM_FACTORY_FORMAT, MOCK_METADATA_SYSTEM), MockSystemFactory.class.getName()); + + // add some serde configs for the in-memory system + mapConfig.putAll(baseInMemorySystemConfigs()); mapConfig.putAll(config); SamzaSqlApplicationRunner runner = new SamzaSqlApplicationRunner(true, new MapConfig(mapConfig)); 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 0cec337700..2f18ab2199 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 @@ -33,7 +33,6 @@ import org.apache.avro.generic.GenericRecord; import org.apache.calcite.plan.RelOptUtil; import org.apache.samza.config.MapConfig; -import org.apache.samza.serializers.JsonSerdeV2Factory; import org.apache.samza.sql.runner.SamzaSqlApplicationConfig; import org.apache.samza.sql.system.TestAvroSystemFactory; import org.apache.samza.sql.util.JsonUtil; @@ -56,17 +55,7 @@ public class TestSamzaSqlEndToEnd extends SamzaSqlIntegrationTestHarness { @Before public void setUp() { - super.setUp(); - configs.put("systems.kafka.samza.factory", "org.apache.samza.system.kafka.KafkaSystemFactory"); - configs.put("systems.kafka.producer.bootstrap.servers", bootstrapUrl()); - configs.put("systems.kafka.consumer.zookeeper.connect", zkConnect()); - configs.put("systems.kafka.samza.key.serde", "object"); - configs.put("systems.kafka.samza.msg.serde", "samzaSqlRelMsg"); - configs.put("systems.kafka.default.stream.replication.factor", "1"); - configs.put("job.default.system", "kafka"); - - configs.put("serializers.registry.object.class", JsonSerdeV2Factory.class.getName()); - configs.put("serializers.registry.samzaSqlRelMsg.class", JsonSerdeV2Factory.class.getName()); + this.configs.clear(); } @Test From 093dc55644d50e4822639b48cc8b01063914ffcf Mon Sep 17 00:00:00 2001 From: Cameron Lee Date: Wed, 21 Aug 2019 09:34:40 -0700 Subject: [PATCH 2/2] update comment, remove unused configs var --- .../SamzaSqlIntegrationTestHarness.java | 2 +- .../test/samzasql/TestSamzaSqlEndToEnd.java | 92 +++++++++---------- 2 files changed, 42 insertions(+), 52 deletions(-) diff --git a/samza-test/src/test/java/org/apache/samza/test/samzasql/SamzaSqlIntegrationTestHarness.java b/samza-test/src/test/java/org/apache/samza/test/samzasql/SamzaSqlIntegrationTestHarness.java index d32d82e6c1..28bad3c083 100644 --- a/samza-test/src/test/java/org/apache/samza/test/samzasql/SamzaSqlIntegrationTestHarness.java +++ b/samza-test/src/test/java/org/apache/samza/test/samzasql/SamzaSqlIntegrationTestHarness.java @@ -46,7 +46,7 @@ protected void runApplication(Config config) { mapConfig.put(JobConfig.JOB_COORDINATOR_SYSTEM, MOCK_METADATA_SYSTEM); mapConfig.put(String.format(SystemConfig.SYSTEM_FACTORY_FORMAT, MOCK_METADATA_SYSTEM), MockSystemFactory.class.getName()); - // add some serde configs for the in-memory system + // configs for using in-memory system as the default system mapConfig.putAll(baseInMemorySystemConfigs()); mapConfig.putAll(config); 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 2f18ab2199..ac84fe8e24 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 @@ -41,7 +41,6 @@ import org.apache.samza.sql.util.SamzaSqlTestConfig; import org.apache.samza.system.OutgoingMessageEnvelope; import org.junit.Assert; -import org.junit.Before; import org.junit.Ignore; import org.junit.Test; import org.slf4j.Logger; @@ -49,21 +48,14 @@ public class TestSamzaSqlEndToEnd extends SamzaSqlIntegrationTestHarness { - private static final Logger LOG = LoggerFactory.getLogger(TestSamzaSqlEndToEnd.class); - private final Map configs = new HashMap<>(); - - @Before - public void setUp() { - this.configs.clear(); - } @Test public void testEndToEnd() { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql = "Insert into testavro.simpleOutputTopic select * from testavro.SIMPLE1"; List sqlStmts = Arrays.asList(sql); staticConfigs.put(SamzaSqlApplicationConfig.CFG_SQL_STMTS_JSON, JsonUtil.toJson(sqlStmts)); @@ -82,7 +74,7 @@ public void testEndToEndWithSystemMessages() { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String avroSamzaToRelMsgConverterDomain = String.format(SamzaSqlApplicationConfig.CFG_FMT_SAMZA_REL_CONVERTER_DOMAIN, "avro"); staticConfigs.put(avroSamzaToRelMsgConverterDomain + SamzaSqlApplicationConfig.CFG_FACTORY, @@ -104,7 +96,7 @@ public void testEndToEndDisableSystemMessages() { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String avroSamzaToRelMsgConverterDomain = String.format(SamzaSqlApplicationConfig.CFG_FMT_SAMZA_REL_CONVERTER_DOMAIN, "avro"); staticConfigs.put(avroSamzaToRelMsgConverterDomain + SamzaSqlApplicationConfig.CFG_FACTORY, @@ -128,7 +120,7 @@ public void testEndToEndWithNullRecords() { TestAvroSystemFactory.messages.clear(); Map staticConfigs = - SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages, false, true); + SamzaSqlTestConfig.fetchStaticConfigsWithFactories(Collections.emptyMap(), numMessages, false, true); String sql = "Insert into testavro.simpleOutputTopic select * from testavro.SIMPLE1"; List sqlStmts = Arrays.asList(sql); staticConfigs.put(SamzaSqlApplicationConfig.CFG_SQL_STMTS_JSON, JsonUtil.toJson(sqlStmts)); @@ -153,7 +145,7 @@ public void testEndToEndWithDifferentSystemSameStream() { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql = "Insert into testavro2.SIMPLE1 select * from testavro.SIMPLE1"; List sqlStmts = Arrays.asList(sql); staticConfigs.put(SamzaSqlApplicationConfig.CFG_SQL_STMTS_JSON, JsonUtil.toJson(sqlStmts)); @@ -171,7 +163,7 @@ public void testEndToEndWithDifferentSystemSameStream() { public void testEndToEndMultiSqlStmts() { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.simpleOutputTopic select * from testavro.SIMPLE1"; String sql2 = "Insert into testavro.SIMPLE3 select * from testavro.SIMPLE2"; List sqlStmts = Arrays.asList(sql1, sql2); @@ -191,7 +183,7 @@ public void testEndToEndMultiSqlStmts() { public void testEndToEndMultiSqlStmtsWithSameSystemStreamAsInputAndOutput() { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.SIMPLE1 select * from testavro.SIMPLE2"; String sql2 = "Insert into testavro.simpleOutputTopic select * from testavro.SIMPLE1"; List sqlStmts = Arrays.asList(sql1, sql2); @@ -212,7 +204,7 @@ public void testEndToEndMultiSqlStmtsWithSameSystemStreamAsInputAndOutput() { public void testEndToEndFanIn() { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.simpleOutputTopic select * from testavro.SIMPLE2"; String sql2 = "Insert into testavro.simpleOutputTopic select * from testavro.SIMPLE1"; List sqlStmts = Arrays.asList(sql1, sql2); @@ -232,7 +224,7 @@ public void testEndToEndFanIn() { public void testEndToEndFanOut() { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.SIMPLE2 select * from testavro.SIMPLE1"; String sql2 = "Insert into testavro.SIMPLE3 select * from testavro.SIMPLE1"; List sqlStmts = Arrays.asList(sql1, sql2); @@ -253,7 +245,7 @@ public void testEndToEndWithProjection() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic(id, long_value) " + " select id, TIMESTAMPDIFF(HOUR, CURRENT_TIMESTAMP(), LOCALTIMESTAMP()) + MONTH(CURRENT_DATE()) as long_value from testavro.SIMPLE1"; List sqlStmts = Arrays.asList(sql1); @@ -273,7 +265,7 @@ public void testEndToEndWithBooleanCheck() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic" + " select * from testavro.COMPLEX1 where bool_value IS TRUE"; List sqlStmts = Arrays.asList(sql1); @@ -290,7 +282,7 @@ public void testEndToEndCompoundBooleanCheck() { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic" + " select * from testavro.COMPLEX1 where id >= 0 and bool_value IS TRUE"; List sqlStmts = Arrays.asList(sql1); @@ -307,7 +299,7 @@ public void testEndToEndCompoundBooleanCheckWorkaround() { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); // BUG Compound boolean checks dont work in calcite, So workaround by casting it to String String sql1 = "Insert into testavro.outputTopic" + " select * from testavro.COMPLEX1 where id >= 0 and CAST(bool_value AS VARCHAR) = 'TRUE'"; @@ -325,7 +317,7 @@ public void testEndToEndWithProjectionWithCase() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic(id, long_value) " + " select id, NOT(id = 5) as bool_value, CASE WHEN id IN (5, 6, 7) THEN CAST('foo' AS VARCHAR) WHEN id < 5 THEN CAST('bars' AS VARCHAR) ELSE NULL END as string_value from testavro.SIMPLE1"; List sqlStmts = Arrays.asList(sql1); @@ -345,7 +337,7 @@ public void testEndToEndWithLike() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic(id, long_value) " + " select id, name as string_value from testavro.SIMPLE1 where name like 'Name%'"; List sqlStmts = Arrays.asList(sql1); @@ -364,7 +356,7 @@ public void testEndToEndWithLike() throws Exception { public void testEndToEndFlatten() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); LOG.info(" Class Path : " + RelOptUtil.class.getProtectionDomain().getCodeSource().getLocation().toURI().getPath()); String sql1 = @@ -390,7 +382,7 @@ public void testEndToEndFlatten() throws Exception { public void testEndToEndComplexRecord() { int numMessages = 10; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic" @@ -410,7 +402,7 @@ public void testEndToEndComplexRecord() { public void testEndToEndNestedRecord() { int numMessages = 10; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic" @@ -429,7 +421,7 @@ public void testEndToEndNestedRecord() { public void testEndToEndFlattenWithUdf() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic(id) select Flatten(MyTestArray(id)) as id from testavro.SIMPLE1"; List sqlStmts = Collections.singletonList(sql1); @@ -450,7 +442,7 @@ public void testEndToEndFlattenWithUdf() throws Exception { public void testEndToEndSubQuery() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic(id) select Flatten(a) as id from (select MyTestArray(id) a from testavro.SIMPLE1)"; List sqlStmts = Collections.singletonList(sql1); @@ -471,7 +463,7 @@ public void testEndToEndSubQuery() throws Exception { public void testUdfUnTypedArgumentToTypedUdf() { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic(id, long_value) " + "select id, MyTest(MyTestObj(id)) as long_value from testavro.SIMPLE1"; List sqlStmts = Collections.singletonList(sql1); @@ -491,7 +483,7 @@ public void testUdfUnTypedArgumentToTypedUdf() { public void testEndToEndUdf() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic(id, long_value) " + "select id, MYTest(id) as long_value from testavro.SIMPLE1"; List sqlStmts = Collections.singletonList(sql1); @@ -515,7 +507,7 @@ public void testEndToEndUdf() throws Exception { public void testEndToEndUdfWithDisabledArgCheck() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.PROFILE1(id, address) " + "select id, BuildOutputRecord('key', GetNestedField(address, 'zip')) as address from testavro.PROFILE"; List sqlStmts = Collections.singletonList(sql1); @@ -535,7 +527,7 @@ public void testEndToEndUdfWithDisabledArgCheck() throws Exception { public void testEndToEndUdfPolymorphism() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic(id, long_value) " + "select MyTestPoly(id) as long_value, MyTestPoly(name) as id from testavro.SIMPLE1"; List sqlStmts = Collections.singletonList(sql1); @@ -559,7 +551,7 @@ public void testEndToEndUdfPolymorphism() throws Exception { public void testRegexMatchUdfInWhereClause() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic(id) " + "select id " @@ -579,8 +571,7 @@ public void testEndToEndStreamTableInnerJoin() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); - staticConfigs.putAll(configs); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql = "Insert into testavro.enrichedPageViewTopic " + "select pv.pageKey as __key__, pv.pageKey as pageKey, coalesce(null, 'N/A') as companyName," @@ -608,8 +599,7 @@ public void testEndToEndStreamTableInnerJoinWithPrimaryKey() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); - staticConfigs.putAll(configs); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql = "Insert into testavro.enrichedPageViewTopic " + "select pv.pageKey as __key__, pv.pageKey as pageKey, coalesce(null, 'N/A') as companyName," @@ -637,8 +627,7 @@ public void testEndToEndStreamTableInnerJoinWithUdf() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); - staticConfigs.putAll(configs); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql = "Insert into testavro.enrichedPageViewTopic " + "select pv.pageKey as __key__, pv.pageKey as pageKey, coalesce(null, 'N/A') as companyName," @@ -666,8 +655,7 @@ public void testEndToEndStreamTableInnerJoinWithNestedRecord() throws Exception int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); - staticConfigs.putAll(configs); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql = "Insert into testavro.enrichedPageViewTopic " + "select pv.pageKey as __key__, pv.pageKey as pageKey, p.name as companyName, p.name as profileName," @@ -700,8 +688,7 @@ public void testEndToEndStreamTableInnerJoinWithFilter() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); - staticConfigs.putAll(configs); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql = "Insert into testavro.enrichedPageViewTopic " + "select pv.pageKey as __key__, pv.pageKey as pageKey, p.name as companyName, p.name as profileName," @@ -734,7 +721,8 @@ public void testEndToEndStreamTableInnerJoinWithNullForeignKeys() throws Excepti int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages, true); + Map staticConfigs = + SamzaSqlTestConfig.fetchStaticConfigsWithFactories(Collections.emptyMap(), numMessages, true); String sql = "Insert into testavro.enrichedPageViewTopic " + "select pv.pageKey as __key__, pv.pageKey as pageKey, p.name as companyName, p.name as profileName," @@ -763,7 +751,8 @@ public void testEndToEndStreamTableLeftJoin() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages, true); + Map staticConfigs = + SamzaSqlTestConfig.fetchStaticConfigsWithFactories(Collections.emptyMap(), numMessages, true); String sql = "Insert into testavro.enrichedPageViewTopic " + "select pv.pageKey as __key__, pv.pageKey as pageKey, p.name as companyName, p.name as profileName," @@ -792,7 +781,8 @@ public void testEndToEndStreamTableRightJoin() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages, true); + Map staticConfigs = + SamzaSqlTestConfig.fetchStaticConfigsWithFactories(Collections.emptyMap(), numMessages, true); String sql = "Insert into testavro.enrichedPageViewTopic " + "select pv.pageKey as __key__, pv.pageKey as pageKey, p.name as companyName, p.name as profileName," @@ -822,7 +812,7 @@ public void testEndToEndStreamTableTableJoin() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql = "Insert into testavro.enrichedPageViewTopic " + "select pv.pageKey as __key__, pv.pageKey as pageKey, c.name as companyName, p.name as profileName," @@ -852,7 +842,7 @@ public void testEndToEndStreamTableTableJoinWithPrimaryKeys() throws Exception { int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql = "Insert into testavro.enrichedPageViewTopic " + "select pv.pageKey as __key__, pv.pageKey as pageKey, c.name as companyName, p.name as profileName," @@ -882,7 +872,7 @@ public void testEndToEndStreamTableTableJoinWithCompositeKey() throws Exception int numMessages = 20; TestAvroSystemFactory.messages.clear(); - Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages); + Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql = "Insert into testavro.enrichedPageViewTopic " + "select pv.pageKey as __key__, pv.pageKey as pageKey, c.name as companyName, p.name as profileName," @@ -917,8 +907,8 @@ public void testEndToEndGroupBy() throws Exception { TestAvroSystemFactory.messages.clear(); Map staticConfigs = - SamzaSqlTestConfig.fetchStaticConfigsWithFactories(configs, numMessages, false, false, windowDurationMs); - staticConfigs.putAll(configs); + SamzaSqlTestConfig.fetchStaticConfigsWithFactories(Collections.emptyMap(), numMessages, false, false, + windowDurationMs); String sql = "Insert into testavro.pageViewCountTopic" + " select 'SampleJob' as jobName, pv.pageKey, count(*) as `count`"