We read every piece of feedback, and take your input very seriously.
To see all available qualifiers, see our documentation.
There was an error while loading. Please reload this page.
1 parent c01fbe1 commit 05c7bafCopy full SHA for 05c7baf
rocketmq-streams-commons/src/main/java/org/apache/rocketmq/streams/common/channel/AbstractChannel.java
@@ -93,6 +93,16 @@ protected void setJsonObject(JSONObject jsonObject) {
93
jsonObject.put("source", Base64Utils.encode(InstantiationUtil.serializeObject(source)));
94
}
95
96
+ @Override
97
+ public void removeSplit(Set<String> splitIds) {
98
+ source.removeSplit(splitIds);
99
+ }
100
+
101
102
+ public void addNewSplit(Set<String> splitIds) {
103
+ source.addNewSplit(splitIds);
104
105
106
@Override
107
public Map<String, MessageOffset> getFinishedQueueIdAndOffsets(CheckPointMessage checkPointMessage) {
108
return sink.getFinishedQueueIdAndOffsets(checkPointMessage);
0 commit comments