diff --git a/core/src/main/java/org/apache/rocketmq/streams/core/rstream/RStream.java b/core/src/main/java/org/apache/rocketmq/streams/core/rstream/RStream.java index 7ff181be..85a6fb73 100644 --- a/core/src/main/java/org/apache/rocketmq/streams/core/rstream/RStream.java +++ b/core/src/main/java/org/apache/rocketmq/streams/core/rstream/RStream.java @@ -27,7 +27,7 @@ public interface RStream { RStream map(ValueMapperAction mapperAction); - RStream flatMap(final ValueMapperAction> mapper); + RStream flatMap(final ValueMapperAction> mapper); RStream filter(FilterAction predictor); diff --git a/core/src/main/java/org/apache/rocketmq/streams/core/rstream/RStreamImpl.java b/core/src/main/java/org/apache/rocketmq/streams/core/rstream/RStreamImpl.java index f9da6ee7..d84626f1 100644 --- a/core/src/main/java/org/apache/rocketmq/streams/core/rstream/RStreamImpl.java +++ b/core/src/main/java/org/apache/rocketmq/streams/core/rstream/RStreamImpl.java @@ -75,7 +75,7 @@ public RStream map(ValueMapperAction mapperAction) { } @Override - public RStream flatMap(ValueMapperAction> mapper) { + public RStream flatMap(ValueMapperAction> mapper) { String name = OperatorNameMaker.makeName(FLAT_MAP_PREFIX, pipeline.getJobId()); MultiValueChangeSupplier changeSupplier = new MultiValueChangeSupplier<>(mapper);