From 83e1745248503bc57bc2a54152d6a5fc7a5b4f41 Mon Sep 17 00:00:00 2001 From: John Roesler Date: Fri, 21 Sep 2018 11:18:07 -0500 Subject: [PATCH] MINOR: rename InternalProcessorContext.initialized --- .../streams/processor/internals/AbstractProcessorContext.java | 2 +- .../streams/processor/internals/GlobalStateUpdateTask.java | 2 +- .../streams/processor/internals/InternalProcessorContext.java | 2 +- .../apache/kafka/streams/processor/internals/StandbyTask.java | 2 +- .../apache/kafka/streams/processor/internals/StreamTask.java | 2 +- .../processor/internals/AbstractProcessorContextTest.java | 2 +- .../org/apache/kafka/test/InternalMockProcessorContext.java | 2 +- .../test/java/org/apache/kafka/test/NoOpProcessorContext.java | 2 +- 8 files changed, 8 insertions(+), 8 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractProcessorContext.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractProcessorContext.java index 0c3fcf201462e..0753b2a8e966d 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractProcessorContext.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractProcessorContext.java @@ -203,7 +203,7 @@ public ThreadCache getCache() { } @Override - public void initialized() { + public void initialize() { initialized = true; } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStateUpdateTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStateUpdateTask.java index d38771362ed36..202fa36b5a9b4 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStateUpdateTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStateUpdateTask.java @@ -74,7 +74,7 @@ public Map initialize() { ); } initTopology(); - processorContext.initialized(); + processorContext.initialize(); return stateMgr.checkpointed(); } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalProcessorContext.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalProcessorContext.java index 6db7a3dc64079..98511fd53dec2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalProcessorContext.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalProcessorContext.java @@ -57,7 +57,7 @@ public interface InternalProcessorContext extends ProcessorContext { /** * Mark this context as being initialized */ - void initialized(); + void initialize(); /** * Mark this context as being uninitialized diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java index 3ac64146187cf..6f4e61787dee4 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java @@ -69,7 +69,7 @@ public boolean initializeStateStores() { log.trace("Initializing state stores"); registerStateStores(); checkpointedOffsets = Collections.unmodifiableMap(stateMgr.checkpointed()); - processorContext.initialized(); + processorContext.initialize(); taskInitialized = true; return true; } 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 2f97b7f27ba29..2ad0acc89d4c1 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 @@ -274,7 +274,7 @@ public void initializeTopology() { transactionInFlight = true; } - processorContext.initialized(); + processorContext.initialize(); taskInitialized = true; diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/AbstractProcessorContextTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/AbstractProcessorContextTest.java index 2869826a421f7..070dba8efa06c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/AbstractProcessorContextTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/AbstractProcessorContextTest.java @@ -58,7 +58,7 @@ public void before() { @Test public void shouldThrowIllegalStateExceptionOnRegisterWhenContextIsInitialized() { - context.initialized(); + context.initialize(); try { context.register(stateStore, null); fail("should throw illegal state exception when context already initialized"); diff --git a/streams/src/test/java/org/apache/kafka/test/InternalMockProcessorContext.java b/streams/src/test/java/org/apache/kafka/test/InternalMockProcessorContext.java index 6d4f5e25419f5..bedf8ebf8c497 100644 --- a/streams/src/test/java/org/apache/kafka/test/InternalMockProcessorContext.java +++ b/streams/src/test/java/org/apache/kafka/test/InternalMockProcessorContext.java @@ -170,7 +170,7 @@ public Serde valueSerde() { // state mgr will be overridden by the state dir and store maps @Override - public void initialized() {} + public void initialize() {} public void setStreamTime(final long currentTime) { streamTime = currentTime; diff --git a/streams/src/test/java/org/apache/kafka/test/NoOpProcessorContext.java b/streams/src/test/java/org/apache/kafka/test/NoOpProcessorContext.java index 42ee8fb33ac11..ce9838919f0b8 100644 --- a/streams/src/test/java/org/apache/kafka/test/NoOpProcessorContext.java +++ b/streams/src/test/java/org/apache/kafka/test/NoOpProcessorContext.java @@ -81,7 +81,7 @@ public void commit() { } @Override - public void initialized() { + public void initialize() { initialized = true; }