diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ResponseJoinProcessorSupplier.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ResponseJoinProcessorSupplier.java index 2fffa89b348e5..600b28078b9fc 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ResponseJoinProcessorSupplier.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ResponseJoinProcessorSupplier.java @@ -29,6 +29,8 @@ import org.apache.kafka.streams.processor.api.Record; import org.apache.kafka.streams.state.ValueAndTimestamp; import org.apache.kafka.streams.state.internals.Murmur3; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.util.function.Supplier; @@ -43,6 +45,7 @@ * @param Type of joined result of primary and foreign values */ public class ResponseJoinProcessorSupplier implements ProcessorSupplier, K, VR> { + private static final Logger LOG = LoggerFactory.getLogger(ResponseJoinProcessorSupplier.class); private final KTableValueGetterSupplier valueGetterSupplier; private final Serializer constructionTimeValueSerializer; private final Supplier valueHashSerdePseudoTopicSupplier; @@ -107,6 +110,8 @@ public void process(final Record> record) { result = joiner.apply(currentValueWithTimestamp == null ? null : currentValueWithTimestamp.value(), record.value().getForeignValue()); } context().forward(record.withValue(result)); + } else { + LOG.trace("Dropping FK-join response due to hash mismatch. Expected {}. Actual {}", messageHash, currentHash); } } };