You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@flink.apache.org by pn...@apache.org on 2019/07/04 11:30:23 UTC
[flink] 05/06: [hotfix][operator] Rename StreamInputProcessor to
StreamOneInputProcessor
This is an automated email from the ASF dual-hosted git repository.
pnowojski pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git
commit de44130b7d275f6e8860c8866fb6cd56caf05175
Author: Piotr Nowojski <pi...@gmail.com>
AuthorDate: Tue Jun 25 15:05:14 2019 +0200
[hotfix][operator] Rename StreamInputProcessor to StreamOneInputProcessor
This is done in order to free `StreamInputProcessor` name as a base interface of
all of the processors.
---
.../org/apache/flink/streaming/runtime/io/InputProcessorUtil.java | 2 +-
.../io/{StreamInputProcessor.java => StreamOneInputProcessor.java} | 6 +++---
.../apache/flink/streaming/runtime/tasks/OneInputStreamTask.java | 6 +++---
3 files changed, 7 insertions(+), 7 deletions(-)
diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/InputProcessorUtil.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/InputProcessorUtil.java
index 4f2a07c..80c0396 100644
--- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/InputProcessorUtil.java
+++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/InputProcessorUtil.java
@@ -34,7 +34,7 @@ import static org.apache.flink.util.Preconditions.checkState;
/**
* Utility for creating {@link CheckpointedInputGate} based on checkpoint mode
- * for {@link StreamInputProcessor} and {@link StreamTwoInputProcessor}.
+ * for {@link StreamOneInputProcessor} and {@link StreamTwoInputProcessor}.
*/
@Internal
public class InputProcessorUtil {
diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/StreamInputProcessor.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/StreamOneInputProcessor.java
similarity index 98%
rename from flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/StreamInputProcessor.java
rename to flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/StreamOneInputProcessor.java
index d6fcad2..b71aaa3 100644
--- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/StreamInputProcessor.java
+++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/StreamOneInputProcessor.java
@@ -62,9 +62,9 @@ import static org.apache.flink.util.Preconditions.checkState;
* @param <IN> The type of the record that can be read with this record reader.
*/
@Internal
-public class StreamInputProcessor<IN> {
+public class StreamOneInputProcessor<IN> {
- private static final Logger LOG = LoggerFactory.getLogger(StreamInputProcessor.class);
+ private static final Logger LOG = LoggerFactory.getLogger(StreamOneInputProcessor.class);
private final StreamTaskInput input;
@@ -87,7 +87,7 @@ public class StreamInputProcessor<IN> {
private Counter numRecordsIn;
@SuppressWarnings("unchecked")
- public StreamInputProcessor(
+ public StreamOneInputProcessor(
InputGate[] inputGates,
TypeSerializer<IN> inputSerializer,
StreamTask<?, ?> checkpointedTask,
diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/tasks/OneInputStreamTask.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/tasks/OneInputStreamTask.java
index e0b328c..c8c51bb 100644
--- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/tasks/OneInputStreamTask.java
+++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/tasks/OneInputStreamTask.java
@@ -26,7 +26,7 @@ import org.apache.flink.runtime.io.network.partition.consumer.InputGate;
import org.apache.flink.runtime.metrics.MetricNames;
import org.apache.flink.streaming.api.graph.StreamConfig;
import org.apache.flink.streaming.api.operators.OneInputStreamOperator;
-import org.apache.flink.streaming.runtime.io.StreamInputProcessor;
+import org.apache.flink.streaming.runtime.io.StreamOneInputProcessor;
import org.apache.flink.streaming.runtime.metrics.WatermarkGauge;
import javax.annotation.Nullable;
@@ -37,7 +37,7 @@ import javax.annotation.Nullable;
@Internal
public class OneInputStreamTask<IN, OUT> extends StreamTask<OUT, OneInputStreamOperator<IN, OUT>> {
- private StreamInputProcessor<IN> inputProcessor;
+ private StreamOneInputProcessor<IN> inputProcessor;
private final WatermarkGauge inputWatermarkGauge = new WatermarkGauge();
@@ -77,7 +77,7 @@ public class OneInputStreamTask<IN, OUT> extends StreamTask<OUT, OneInputStreamO
if (numberOfInputs > 0) {
InputGate[] inputGates = getEnvironment().getAllInputGates();
- inputProcessor = new StreamInputProcessor<>(
+ inputProcessor = new StreamOneInputProcessor<>(
inputGates,
inSerializer,
this,