From 2c26b9fb284b8a1eabef083a9abc970a439deb7e Mon Sep 17 00:00:00 2001 From: Yasuhiro Matsuda Date: Tue, 22 Sep 2015 12:56:32 -0700 Subject: [PATCH] set the initialized flag in context --- .../streams/processor/internals/ProcessorContextImpl.java | 4 ++++ .../apache/kafka/streams/processor/internals/StreamTask.java | 5 +++-- 2 files changed, 7 insertions(+), 2 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java index b01dbc1d05667..b3502220cf3ed 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java @@ -79,6 +79,10 @@ public RecordCollector recordCollector() { return this.collector; } + public void initialized() { + this.initialized = true; + } + @Override public boolean joinable() { Set partitions = this.task.partitions(); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index be155186054de..86dc735e14ae2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -26,7 +26,6 @@ import org.apache.kafka.common.metrics.Metrics; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.streams.StreamingConfig; -import org.apache.kafka.streams.processor.ProcessorContext; import org.apache.kafka.streams.processor.TimestampExtractor; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -53,7 +52,7 @@ public class StreamTask implements Punctuator { private final PartitionGroup partitionGroup; private final PartitionGroup.RecordInfo recordInfo = new PartitionGroup.RecordInfo(); private final PunctuationQueue punctuationQueue; - private final ProcessorContext processorContext; + private final ProcessorContextImpl processorContext; private final ProcessorTopology topology; private final Map consumedOffsets; @@ -135,6 +134,8 @@ public StreamTask(int id, this.currNode = null; } } + + this.processorContext.initialized(); } public int id() {