This is an automated email from the ASF dual-hosted git repository.
ruanwenjun pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/dolphinscheduler.git
The following commit(s) were added to refs/heads/dev by this push:
new f16479955a [Fix-18640][Registry] Preserve previous values in etcd
REMOVE events (#18641)
f16479955a is described below
commit f16479955a28d0d0d51e5ccea5fb975da8db8417
Author: michaellx1057 <[email protected]>
AuthorDate: Wed Sep 23 22:14:35 2026 +0800
[Fix-18640][Registry] Preserve previous values in etcd REMOVE events
(#18641)
---
.../plugin/registry/etcd/EtcdRegistry.java | 4 +-
.../registry/etcd/EtcdRegistryEventTest.java | 87 ++++++++++++++++++++++
.../plugin/registry/RegistryTestCase.java | 52 +++++++++++++
3 files changed, 142 insertions(+), 1 deletion(-)
diff --git
a/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/main/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistry.java
b/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/main/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistry.java
index a5955d3212..048def65e4 100644
---
a/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/main/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistry.java
+++
b/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/main/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistry.java
@@ -438,7 +438,9 @@ public class EtcdRegistry implements Registry {
.watchedPath(watchedPath)
.eventPath(Optional.ofNullable(keyValue).map(kv ->
kv.getKey().toString(StandardCharsets.UTF_8))
.orElse(null))
- .eventData(Optional.ofNullable(keyValue).map(kv ->
kv.getValue().toString(StandardCharsets.UTF_8))
+ .eventData(Optional
+ .ofNullable(eventType == Event.Type.REMOVE ?
watchEvent.getPrevKV() : watchEvent.getKeyValue())
+ .map(kv ->
kv.getValue().toString(StandardCharsets.UTF_8))
.orElse(null))
.build();
}
diff --git
a/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/test/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistryEventTest.java
b/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/test/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistryEventTest.java
new file mode 100644
index 0000000000..c69bb40610
--- /dev/null
+++
b/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/test/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistryEventTest.java
@@ -0,0 +1,87 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.dolphinscheduler.plugin.registry.etcd;
+
+import org.apache.dolphinscheduler.registry.api.Event;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import com.google.protobuf.ByteString;
+
+import io.etcd.jetcd.ByteSequence;
+import io.etcd.jetcd.KeyValue;
+import io.etcd.jetcd.watch.WatchEvent;
+
+class EtcdRegistryEventTest {
+
+ private static final String WATCHED_PATH = "/nodes";
+ private static final String EVENT_PATH = "/nodes/master-1";
+
+ private final EtcdRegistry registry = Mockito.mock(EtcdRegistry.class);
+
+ @Test
+ void testDeleteUsesPreviousValue() {
+ // An etcd DELETE contains only the key and modification revision in
its current KV.
+ KeyValue deletedKeyValue = new
KeyValue(io.etcd.jetcd.api.KeyValue.newBuilder()
+ .setKey(ByteString.copyFromUtf8(EVENT_PATH))
+ .setModRevision(3)
+ .build(), ByteSequence.EMPTY);
+ assertEvent(new WatchEvent(deletedKeyValue, keyValue(EVENT_PATH,
"previous-heartbeat"),
+ WatchEvent.EventType.DELETE), Event.Type.REMOVE,
"previous-heartbeat");
+ }
+
+ @Test
+ void testAddUsesCurrentValue() {
+ assertEvent(new WatchEvent(keyValue(EVENT_PATH, "current-heartbeat"),
+ new KeyValue(io.etcd.jetcd.api.KeyValue.getDefaultInstance(),
ByteSequence.EMPTY),
+ WatchEvent.EventType.PUT), Event.Type.ADD,
"current-heartbeat");
+ }
+
+ @Test
+ void testUpdateUsesCurrentValue() {
+ assertEvent(new WatchEvent(keyValue(EVENT_PATH, "current-heartbeat"),
+ keyValue(EVENT_PATH, "previous-heartbeat"),
WatchEvent.EventType.PUT),
+ Event.Type.UPDATE, "current-heartbeat");
+ }
+
+ @Test
+ void testDeleteWithoutPreviousValuePreservesPath() {
+ assertEvent(new WatchEvent(keyValue(EVENT_PATH, ""),
+ new KeyValue(io.etcd.jetcd.api.KeyValue.getDefaultInstance(),
ByteSequence.EMPTY),
+ WatchEvent.EventType.DELETE), Event.Type.REMOVE, "");
+ }
+
+ private void assertEvent(WatchEvent watchEvent, Event.Type expectedType,
String expectedData) {
+ Event event = ReflectionTestUtils.invokeMethod(registry, "toEvent",
watchEvent, WATCHED_PATH);
+ Assertions.assertNotNull(event);
+ Assertions.assertEquals(expectedType, event.getType());
+ Assertions.assertEquals(WATCHED_PATH, event.getWatchedPath());
+ Assertions.assertEquals(EVENT_PATH, event.getEventPath());
+ Assertions.assertEquals(expectedData, event.getEventData());
+ }
+
+ private KeyValue keyValue(String key, String value) {
+ return new KeyValue(io.etcd.jetcd.api.KeyValue.newBuilder()
+ .setKey(ByteString.copyFromUtf8(key))
+ .setValue(ByteString.copyFromUtf8(value))
+ .build(), ByteSequence.EMPTY);
+ }
+}
diff --git
a/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-it/src/test/java/org/apache/dolphinscheduler/plugin/registry/RegistryTestCase.java
b/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-it/src/test/java/org/apache/dolphinscheduler/plugin/registry/RegistryTestCase.java
index b819ef0eee..8e192c25bf 100644
---
a/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-it/src/test/java/org/apache/dolphinscheduler/plugin/registry/RegistryTestCase.java
+++
b/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-it/src/test/java/org/apache/dolphinscheduler/plugin/registry/RegistryTestCase.java
@@ -119,6 +119,58 @@ public abstract class RegistryTestCase<R extends Registry>
{
});
}
+ @SneakyThrows
+ @Test
+ public void testSubscribeEventData() {
+ registry.start();
+
+ // Futures safely publish each first event and its payload to the test
thread.
+ final CompletableFuture<Event> subscribeAdded = new
CompletableFuture<>();
+ final CompletableFuture<Event> subscribeRemoved = new
CompletableFuture<>();
+ final CompletableFuture<Event> subscribeUpdated = new
CompletableFuture<>();
+
+ final SubscribeListener subscribeListener = new SubscribeListener() {
+
+ @Override
+ public void notify(Event event) {
+ // Keep assertions on the test thread so callback error
handling cannot hide failures.
+ if (event.getType() == Event.Type.ADD) {
+ subscribeAdded.complete(event);
+ }
+ if (event.getType() == Event.Type.REMOVE) {
+ subscribeRemoved.complete(event);
+ }
+ if (event.getType() == Event.Type.UPDATE) {
+ subscribeUpdated.complete(event);
+ }
+ }
+
+ @Override
+ public SubscribeScope getSubscribeScope() {
+ return SubscribeScope.PATH_ONLY;
+ }
+ };
+ String key = "/nodes/master" + System.nanoTime();
+ registry.subscribe(key, subscribeListener);
+ // Wait after each change so polling registries cannot collapse
consecutive operations.
+ registry.put(key, "v1", true);
+ assertSubscribeEvent(subscribeAdded.get(10, TimeUnit.SECONDS),
Event.Type.ADD, key, "v1");
+
+ registry.put(key, "v2", true);
+ assertSubscribeEvent(subscribeUpdated.get(10, TimeUnit.SECONDS),
Event.Type.UPDATE, key, "v2");
+
+ // REMOVE must retain the last value before deletion, not the initial
value.
+ registry.delete(key);
+ assertSubscribeEvent(subscribeRemoved.get(10, TimeUnit.SECONDS),
Event.Type.REMOVE, key, "v2");
+ }
+
+ private void assertSubscribeEvent(Event event, Event.Type expectedType,
String key, String expectedData) {
+ Assertions.assertEquals(expectedType, event.getType());
+ Assertions.assertEquals(key, event.getWatchedPath());
+ Assertions.assertEquals(key, event.getEventPath());
+ Assertions.assertEquals(expectedData, event.getEventData());
+ }
+
@SneakyThrows
@Test
public void testAddConnectionStateListener() {