Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@
import org.apache.kafka.streams.kstream.internals.suppress.SuppressedInternal;
import org.apache.kafka.streams.processor.ProcessorSupplier;
import org.apache.kafka.streams.state.KeyValueStore;
import org.apache.kafka.streams.state.StoreBuilder;
import org.apache.kafka.streams.state.internals.InMemoryTimeOrderedKeyValueBuffer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
Expand Down Expand Up @@ -465,28 +466,9 @@ private <VO, VR> KTable<K, VR> doJoin(final KTable<K, VO> other,
final boolean rightOuter) {
Objects.requireNonNull(other, "other can't be null");
Objects.requireNonNull(joiner, "joiner can't be null");
final String internalQueryableName = materializedInternal == null ? null : materializedInternal.storeName();
final String joinMergeName = builder.newProcessorName(MERGE_NAME);

return buildJoin(
(AbstractStream<K, VO>) other,
joiner,
leftOuter,
rightOuter,
joinMergeName,
internalQueryableName,
materializedInternal
);
}

private <V1, R> KTable<K, R> buildJoin(final AbstractStream<K, V1> other,
final ValueJoiner<? super V, ? super V1, ? extends R> joiner,
final boolean leftOuter,
final boolean rightOuter,
final String joinMergeName,
final String internalQueryableName,
final MaterializedInternal<K, R, KeyValueStore<Bytes, byte[]>> materializedInternal) {
final Set<String> allSourceNodes = ensureJoinableWith(other);
final String joinMergeName = builder.newProcessorName(MERGE_NAME);
final Set<String> allSourceNodes = ensureJoinableWith((AbstractStream<K, VO>) other);

if (leftOuter) {
enableSendingOldValues();
Expand All @@ -495,57 +477,67 @@ private <V1, R> KTable<K, R> buildJoin(final AbstractStream<K, V1> other,
((KTableImpl) other).enableSendingOldValues();
}

final String joinThisName = builder.newProcessorName(JOINTHIS_NAME);
final String joinOtherName = builder.newProcessorName(JOINOTHER_NAME);


final KTableKTableAbstractJoin<K, R, V, V1> joinThis;
final KTableKTableAbstractJoin<K, R, V1, V> joinOther;
final KTableKTableAbstractJoin<K, VR, V, VO> joinThis;
final KTableKTableAbstractJoin<K, VR, VO, V> joinOther;

if (!leftOuter) { // inner
joinThis = new KTableKTableInnerJoin<>(this, (KTableImpl<K, ?, V1>) other, joiner);
joinOther = new KTableKTableInnerJoin<>((KTableImpl<K, ?, V1>) other, this, reverseJoiner(joiner));
joinThis = new KTableKTableInnerJoin<>(this, (KTableImpl<K, ?, VO>) other, joiner);
joinOther = new KTableKTableInnerJoin<>((KTableImpl<K, ?, VO>) other, this, reverseJoiner(joiner));
} else if (!rightOuter) { // left
joinThis = new KTableKTableLeftJoin<>(this, (KTableImpl<K, ?, V1>) other, joiner);
joinOther = new KTableKTableRightJoin<>((KTableImpl<K, ?, V1>) other, this, reverseJoiner(joiner));
joinThis = new KTableKTableLeftJoin<>(this, (KTableImpl<K, ?, VO>) other, joiner);
joinOther = new KTableKTableRightJoin<>((KTableImpl<K, ?, VO>) other, this, reverseJoiner(joiner));
} else { // outer
joinThis = new KTableKTableOuterJoin<>(this, (KTableImpl<K, ?, V1>) other, joiner);
joinOther = new KTableKTableOuterJoin<>((KTableImpl<K, ?, V1>) other, this, reverseJoiner(joiner));
joinThis = new KTableKTableOuterJoin<>(this, (KTableImpl<K, ?, VO>) other, joiner);
joinOther = new KTableKTableOuterJoin<>((KTableImpl<K, ?, VO>) other, this, reverseJoiner(joiner));
}

final KTableKTableJoinMerger<K, R> joinMerge = new KTableKTableJoinMerger<>(joinThis, joinOther, internalQueryableName);
final String joinThisName = builder.newProcessorName(JOINTHIS_NAME);
final String joinOtherName = builder.newProcessorName(JOINOTHER_NAME);

final ProcessorParameters<K, Change<V>> joinThisProcessorParameters = new ProcessorParameters<>(joinThis, joinThisName);
final ProcessorParameters<K, Change<VO>> joinOtherProcessorParameters = new ProcessorParameters<>(joinOther, joinOtherName);

final KTableKTableJoinNode.KTableKTableJoinNodeBuilder<K, V, V1, R> kTableJoinNodeBuilder = KTableKTableJoinNode.kTableKTableJoinNodeBuilder();
final Serde<K> keySerde;
final Serde<VR> valueSerde;
final String queryableStoreName;
final StoreBuilder<KeyValueStore<K, VR>> storeBuilder;

// only materialize if specified in Materialized
if (materializedInternal != null) {
kTableJoinNodeBuilder.withMaterializedInternal(materializedInternal);
keySerde = materializedInternal.keySerde() != null ? materializedInternal.keySerde() : this.keySerde;
valueSerde = materializedInternal.valueSerde();
queryableStoreName = materializedInternal.storeName();
storeBuilder = new KeyValueStoreMaterializer<>(materializedInternal).materialize();
} else {
keySerde = this.keySerde;
valueSerde = null;
queryableStoreName = null;
storeBuilder = null;
}
kTableJoinNodeBuilder.withNodeName(joinMergeName);

final ProcessorParameters<K, Change<V>> joinThisProcessorParameters = new ProcessorParameters<>(joinThis, joinThisName);
final ProcessorParameters<K, Change<V1>> joinOtherProcessorParameters = new ProcessorParameters<>(joinOther, joinOtherName);
final ProcessorParameters<K, Change<R>> joinMergeProcessorParameters = new ProcessorParameters<>(joinMerge, joinMergeName);

kTableJoinNodeBuilder.withJoinMergeProcessorParameters(joinMergeProcessorParameters)
.withJoinOtherProcessorParameters(joinOtherProcessorParameters)
.withJoinThisProcessorParameters(joinThisProcessorParameters)
.withJoinThisStoreNames(valueGetterSupplier().storeNames())
.withJoinOtherStoreNames(((KTableImpl) other).valueGetterSupplier().storeNames())
.withOtherJoinSideNodeName(((KTableImpl) other).name)
.withThisJoinSideNodeName(name);

final KTableKTableJoinNode<K, V, V1, R> kTableKTableJoinNode = kTableJoinNodeBuilder.build();
final KTableKTableJoinNode<K, V, VO, VR> kTableKTableJoinNode =
KTableKTableJoinNode.<K, V, VO, VR>kTableKTableJoinNodeBuilder()
.withNodeName(joinMergeName)
.withJoinThisProcessorParameters(joinThisProcessorParameters)
.withJoinOtherProcessorParameters(joinOtherProcessorParameters)
.withThisJoinSideNodeName(name)
.withOtherJoinSideNodeName(((KTableImpl) other).name)
.withJoinThisStoreNames(valueGetterSupplier().storeNames())
.withJoinOtherStoreNames(((KTableImpl) other).valueGetterSupplier().storeNames())
.withKeySerde(keySerde)
.withValueSerde(valueSerde)
.withQueryableStoreName(queryableStoreName)
.withStoreBuilder(storeBuilder)
.build();
builder.addGraphNode(this.streamsGraphNode, kTableKTableJoinNode);

// we can inherit parent key serde if user do not provide specific overrides
return new KTableImpl<K, Change<R>, R>(
joinMergeName,
materializedInternal != null && materializedInternal.keySerde() != null ? materializedInternal.keySerde() : keySerde,
materializedInternal != null ? materializedInternal.valueSerde() : null,
return new KTableImpl<K, Change<VR>, VR>(
kTableKTableJoinNode.nodeName(),
kTableKTableJoinNode.keySerde(),
kTableKTableJoinNode.valueSerde(),
allSourceNodes,
internalQueryableName,
joinMerge,
kTableKTableJoinNode.queryableStoreName(),
kTableKTableJoinNode.joinMerger(),
kTableKTableJoinNode,
builder
);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
import java.util.HashSet;
import java.util.Set;

class KTableKTableJoinMerger<K, V> implements KTableProcessorSupplier<K, V, V> {
public class KTableKTableJoinMerger<K, V> implements KTableProcessorSupplier<K, V, V> {

private final KTableProcessorSupplier<K, ?, V> parent1;
private final KTableProcessorSupplier<K, ?, V> parent2;
Expand All @@ -40,6 +40,10 @@ class KTableKTableJoinMerger<K, V> implements KTableProcessorSupplier<K, V, V> {
this.queryableName = queryableName;
}

public String getQueryableName() {
return queryableName;
}

@Override
public Processor<K, Change<V>> get() {
return new KTableKTableJoinMergeProcessor();
Expand Down Expand Up @@ -78,6 +82,17 @@ public void enableSendingOldValues() {
sendOldValues = true;
}

public static <K, V> KTableKTableJoinMerger<K, V> of(final KTableProcessorSupplier<K, ?, V> parent1,
final KTableProcessorSupplier<K, ?, V> parent2) {
return of(parent1, parent2, null);
}

public static <K, V> KTableKTableJoinMerger<K, V> of(final KTableProcessorSupplier<K, ?, V> parent1,
final KTableProcessorSupplier<K, ?, V> parent2,
final String queryableName) {
return new KTableKTableJoinMerger<>(parent1, parent2, queryableName);
}

private class KTableKTableJoinMergeProcessor extends AbstractProcessor<K, Change<V>> {
private KeyValueStore<K, V> store;
private TupleForwarder<K, V> tupleForwarder;
Expand Down
Loading