diff --git a/ksqldb-engine/src/main/java/io/confluent/ksql/engine/QueryPlan.java b/ksqldb-engine/src/main/java/io/confluent/ksql/engine/QueryPlan.java index a89cade265ed..ff4560ada6e6 100644 --- a/ksqldb-engine/src/main/java/io/confluent/ksql/engine/QueryPlan.java +++ b/ksqldb-engine/src/main/java/io/confluent/ksql/engine/QueryPlan.java @@ -85,12 +85,13 @@ public boolean equals(final Object o) { return Objects.equals(sources, queryPlan.sources) && Objects.equals(sink, queryPlan.sink) && Objects.equals(physicalPlan, queryPlan.physicalPlan) + && Objects.equals(runtimeId, queryPlan.runtimeId) && Objects.equals(queryId, queryPlan.queryId); } @Override public int hashCode() { - return Objects.hash(sources, sink, physicalPlan, queryId); + return Objects.hash(sources, sink, physicalPlan, queryId, runtimeId); } } diff --git a/ksqldb-engine/src/main/java/io/confluent/ksql/query/QueryBuilder.java b/ksqldb-engine/src/main/java/io/confluent/ksql/query/QueryBuilder.java index 0cf3cd48255c..8b90d26b2c8c 100644 --- a/ksqldb-engine/src/main/java/io/confluent/ksql/query/QueryBuilder.java +++ b/ksqldb-engine/src/main/java/io/confluent/ksql/query/QueryBuilder.java @@ -89,8 +89,9 @@ import org.apache.kafka.streams.Topology; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.KTable; +import org.apache.kafka.streams.processor.internals.namedtopology.KafkaStreamsNamedTopologyWrapper; import org.apache.kafka.streams.processor.internals.namedtopology.NamedTopology; -import org.apache.kafka.streams.processor.internals.namedtopology.NamedTopologyStreamsBuilder; +import org.apache.kafka.streams.processor.internals.namedtopology.NamedTopologyBuilder; /** * A builder for creating queries metadata. @@ -380,12 +381,8 @@ PersistentQueryMetadata buildPersistentQueryInSharedRuntime( final SharedKafkaStreamsRuntime sharedKafkaStreamsRuntime = getKafkaStreamsInstance( sources.stream().map(DataSource::getName).collect(Collectors.toSet()), queryId); - - final String applicationId = sharedKafkaStreamsRuntime - .getStreamProperties() - .get(StreamsConfig.APPLICATION_ID_CONFIG) - .toString(); - final Map streamsProperties = sharedKafkaStreamsRuntime.getStreamProperties(); + final String applicationId = sharedKafkaStreamsRuntime.getApplicationId(); + final Map queryOverrides = sharedKafkaStreamsRuntime.getStreamProperties(); final LogicalSchema logicalSchema; final KeyFormat keyFormat; @@ -417,18 +414,21 @@ PersistentQueryMetadata buildPersistentQueryInSharedRuntime( keyFormat.getFeatures(), valueFormat.getFeatures() ); - final NamedTopologyStreamsBuilder namedTopologyStreamsBuilder = new NamedTopologyStreamsBuilder( - queryId.toString() - ); + + final NamedTopologyBuilder namedTopologyBuilder = + ((KafkaStreamsNamedTopologyWrapper) sharedKafkaStreamsRuntime.getKafkaStreams()) + .newNamedTopologyBuilder( + queryId.toString(), + PropertiesUtil.asProperties(queryOverrides) + ); final RuntimeBuildContext runtimeBuildContext = buildContext( applicationId, queryId, - namedTopologyStreamsBuilder + namedTopologyBuilder ); final Object result = buildQueryImplementation(physicalPlan, runtimeBuildContext); - final NamedTopology topology = namedTopologyStreamsBuilder - .buildNamedTopology(PropertiesUtil.asProperties(streamsProperties)); + final NamedTopology topology = namedTopologyBuilder.build(); final Optional materializationProviderBuilder = getMaterializationInfo(result).map(info -> @@ -436,7 +436,7 @@ PersistentQueryMetadata buildPersistentQueryInSharedRuntime( info, querySchema, keyFormat, - streamsProperties, + queryOverrides, applicationId )); @@ -450,7 +450,7 @@ PersistentQueryMetadata buildPersistentQueryInSharedRuntime( querySchema.logicalSchema(), result, allPersistentQueries, - streamsProperties, + queryOverrides, applicationId, ksqlConfig, ksqlTopic, @@ -475,7 +475,7 @@ PersistentQueryMetadata buildPersistentQueryInSharedRuntime( sinkDataSource, listener, classifier, - streamsProperties, + queryOverrides, scalablePushRegistry ); } diff --git a/ksqldb-engine/src/main/java/io/confluent/ksql/query/QueryRegistryImpl.java b/ksqldb-engine/src/main/java/io/confluent/ksql/query/QueryRegistryImpl.java index 4e4f56005351..0ee8a2a17e50 100644 --- a/ksqldb-engine/src/main/java/io/confluent/ksql/query/QueryRegistryImpl.java +++ b/ksqldb-engine/src/main/java/io/confluent/ksql/query/QueryRegistryImpl.java @@ -63,13 +63,13 @@ public class QueryRegistryImpl implements QueryRegistry { private static final BiPredicate FILTER_QUERIES_WITH_SINK = (sourceName, query) -> query.getSinkName().equals(Optional.of(sourceName)); - private final Map persistentQueries; - private final Map allLiveQueries; - private final Map createAsQueries; - private final Map> insertQueries; + private final Map persistentQueries = new ConcurrentHashMap<>(); + private final Map allLiveQueries = new ConcurrentHashMap<>(); + private final Map createAsQueries = new ConcurrentHashMap<>(); + private final Map> insertQueries = new ConcurrentHashMap<>(); private final Collection eventListeners; private final QueryBuilderFactory queryBuilderFactory; - private final List streams; + private final List streams = new ArrayList<>(); private final boolean sandbox; public QueryRegistryImpl(final Collection eventListeners) { @@ -80,23 +80,14 @@ public QueryRegistryImpl(final Collection eventListeners) { final Collection eventListeners, final QueryBuilderFactory queryBuilderFactory ) { - this.persistentQueries = new ConcurrentHashMap<>(); - this.allLiveQueries = new ConcurrentHashMap<>(); - this.createAsQueries = new ConcurrentHashMap<>(); - this.insertQueries = new ConcurrentHashMap<>(); this.eventListeners = Objects.requireNonNull(eventListeners); this.queryBuilderFactory = Objects.requireNonNull(queryBuilderFactory); - this.streams = new ArrayList<>(); this.sandbox = false; } // Used to construct a sandbox private QueryRegistryImpl(final QueryRegistryImpl original) { queryBuilderFactory = original.queryBuilderFactory; - persistentQueries = new ConcurrentHashMap<>(); - allLiveQueries = new ConcurrentHashMap<>(); - createAsQueries = new ConcurrentHashMap<>(); - insertQueries = new ConcurrentHashMap<>(); original.allLiveQueries.forEach((queryId, queryMetadata) -> { if (queryMetadata instanceof PersistentQueryMetadataImpl) { final PersistentQueryMetadata sandboxed = SandboxedPersistentQueryMetadataImpl.of( @@ -132,9 +123,9 @@ private QueryRegistryImpl(final QueryRegistryImpl original) { .filter(Optional::isPresent) .map(Optional::get) .collect(Collectors.toList()); - this.streams = original.streams.stream() + streams.addAll(original.streams.stream() .map(SandboxedSharedKafkaStreamsRuntimeImpl::new) - .collect(Collectors.toList()); + .collect(Collectors.toList())); sandbox = true; } diff --git a/ksqldb-engine/src/main/java/io/confluent/ksql/util/BinPackedPersistentQueryMetadataImpl.java b/ksqldb-engine/src/main/java/io/confluent/ksql/util/BinPackedPersistentQueryMetadataImpl.java index b42e3b2067e3..74c87b3dc1e0 100644 --- a/ksqldb-engine/src/main/java/io/confluent/ksql/util/BinPackedPersistentQueryMetadataImpl.java +++ b/ksqldb-engine/src/main/java/io/confluent/ksql/util/BinPackedPersistentQueryMetadataImpl.java @@ -119,7 +119,7 @@ public BinPackedPersistentQueryMetadataImpl( this.statementString = Objects.requireNonNull(statementString, "statementString"); this.executionPlan = Objects.requireNonNull(executionPlan, "executionPlan"); this.applicationId = Objects.requireNonNull(applicationId, "applicationId"); - this.topology = Objects.requireNonNull(topology, "kafkaTopicClient"); + this.topology = Objects.requireNonNull(topology, "namedTopology"); this.sharedKafkaStreamsRuntime = Objects.requireNonNull(sharedKafkaStreamsRuntime, "sharedKafkaStreamsRuntime"); this.sinkDataSource = requireNonNull(sinkDataSource, "sinkDataSource"); diff --git a/ksqldb-engine/src/test/java/io/confluent/ksql/query/QueryBuilderTest.java b/ksqldb-engine/src/test/java/io/confluent/ksql/query/QueryBuilderTest.java index d2389f0ea9e9..c834386e65fc 100644 --- a/ksqldb-engine/src/test/java/io/confluent/ksql/query/QueryBuilderTest.java +++ b/ksqldb-engine/src/test/java/io/confluent/ksql/query/QueryBuilderTest.java @@ -80,6 +80,9 @@ import org.apache.kafka.streams.Topology; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.processor.internals.namedtopology.KafkaStreamsNamedTopologyWrapper; +import org.apache.kafka.streams.processor.internals.namedtopology.NamedTopology; +import org.apache.kafka.streams.processor.internals.namedtopology.NamedTopologyBuilder; + import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -155,6 +158,8 @@ public class QueryBuilderTest { @Mock private StreamsBuilder streamsBuilder; @Mock + private NamedTopologyBuilder namedTopologyBuilder; + @Mock private FunctionRegistry functionRegistry; @Mock private KafkaStreams kafkaStreams; @@ -163,6 +168,8 @@ public class QueryBuilderTest { @Mock private Topology topology; @Mock + private NamedTopology namedTopology; + @Mock private KsMaterializationFactory ksMaterializationFactory; @Mock private KsMaterialization ksMaterialization; @@ -197,6 +204,7 @@ public void setup() { when(ksqlTopic.getValueFormat()).thenReturn(VALUE_FORMAT); when(kafkaStreamsBuilder.build(any(), any())).thenReturn(kafkaStreams); when(kafkaStreamsBuilder.buildNamedTopologyWrapper(any())).thenReturn(kafkaStreamsNamedTopologyWrapper); + when(kafkaStreamsNamedTopologyWrapper.newNamedTopologyBuilder(any(), any())).thenReturn(namedTopologyBuilder); when(tableHolder.getMaterializationBuilder()).thenReturn(Optional.of(materializationBuilder)); when(materializationBuilder.build()).thenReturn(materializationInfo); when(materializationInfo.getStateStoreSchema()).thenReturn(aggregationSchema); @@ -214,9 +222,9 @@ public void setup() { when(ksqlConfig.getBoolean(KsqlConfig.KSQL_SHARED_RUNTIME_ENABLED)).thenReturn(false); when(physicalPlan.build(any())).thenReturn(tableHolder); when(streamsBuilder.build(any())).thenReturn(topology); + when(namedTopologyBuilder.build()).thenReturn(namedTopology); when(config.getConfig(true)).thenReturn(ksqlConfig); when(config.getOverrides()).thenReturn(OVERRIDES); - when(kstream.filter(any())).thenReturn(kstream); sharedKafkaStreamsRuntimes = new ArrayList<>(); queryBuilder = new QueryBuilder( diff --git a/ksqldb-engine/src/test/java/io/confluent/ksql/util/QueryMetadataTest.java b/ksqldb-engine/src/test/java/io/confluent/ksql/util/QueryMetadataTest.java index c8d06ae95a4c..c2a375757dd5 100644 --- a/ksqldb-engine/src/test/java/io/confluent/ksql/util/QueryMetadataTest.java +++ b/ksqldb-engine/src/test/java/io/confluent/ksql/util/QueryMetadataTest.java @@ -61,8 +61,8 @@ @RunWith(MockitoJUnitRunner.class) public class QueryMetadataTest { - private static long RETRY_BACKOFF_INITIAL_MS = 1; - private static long RETRY_BACKOFF_MAX_MS = 10; + private static final long RETRY_BACKOFF_INITIAL_MS = 1; + private static final long RETRY_BACKOFF_MAX_MS = 10; private static final String QUERY_APPLICATION_ID = "Query1"; private static final QueryId QUERY_ID = new QueryId("queryId"); private static final LogicalSchema SOME_SCHEMA = LogicalSchema.builder()