Skip to content

Commit 05c7baf

Browse files
author
tigerweili
committed
fix code check fail: not ovride removeSplit() and addNewSplit()
1 parent c01fbe1 commit 05c7baf

File tree

1 file changed

+10
-0
lines changed
  • rocketmq-streams-commons/src/main/java/org/apache/rocketmq/streams/common/channel

1 file changed

+10
-0
lines changed

rocketmq-streams-commons/src/main/java/org/apache/rocketmq/streams/common/channel/AbstractChannel.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,16 @@ protected void setJsonObject(JSONObject jsonObject) {
9393
jsonObject.put("source", Base64Utils.encode(InstantiationUtil.serializeObject(source)));
9494
}
9595

96+
@Override
97+
public void removeSplit(Set<String> splitIds) {
98+
source.removeSplit(splitIds);
99+
}
100+
101+
@Override
102+
public void addNewSplit(Set<String> splitIds) {
103+
source.addNewSplit(splitIds);
104+
}
105+
96106
@Override
97107
public Map<String, MessageOffset> getFinishedQueueIdAndOffsets(CheckPointMessage checkPointMessage) {
98108
return sink.getFinishedQueueIdAndOffsets(checkPointMessage);

0 commit comments

Comments
 (0)