Skip to content

Commit 73cf067

Browse files
authored
Pipe: Fixed the bug that air gap receiver may not respond in temporary timeout exception & Optimized the directory check in receiver & Fixed the bug that the "skipIfNoPrivileges" may be wrongly reused at receiver & Optimized the configNode pipe logic (#17556)
* re * sink * fix * rollback
1 parent becc0c4 commit 73cf067

13 files changed

Lines changed: 651 additions & 57 deletions

File tree

iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1666,16 +1666,18 @@ public void pipeHandleLeaderChange(
16661666
}
16671667
}
16681668

1669-
public void pipeHandleMetaChange(
1669+
public boolean pipeHandleMetaChange(
16701670
boolean needWriteConsensusOnConfigNodes, boolean needPushPipeMetaToDataNodes) {
16711671
try {
16721672
final long procedureId =
16731673
executor.submitProcedure(
16741674
new PipeHandleMetaChangeProcedure(
16751675
needWriteConsensusOnConfigNodes, needPushPipeMetaToDataNodes));
16761676
LOGGER.info("PipeHandleMetaChangeProcedure was submitted, procedureId: {}.", procedureId);
1677+
return true;
16771678
} catch (Exception e) {
16781679
LOGGER.warn("PipeHandleMetaChangeProcedure was failed to submit.", e);
1680+
return false;
16791681
}
16801682
}
16811683

iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java

Lines changed: 20 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -59,7 +59,7 @@ public class PipeHeartbeatParser {
5959
this.configManager = configManager;
6060

6161
heartbeatCounter = 0;
62-
registeredNodeNumber = 1;
62+
registeredNodeNumber = getExpectedHeartbeatNodeCount();
6363

6464
needWriteConsensusOnConfigNodes = new AtomicBoolean(false);
6565
needPushPipeMetaToDataNodes = new AtomicBoolean(false);
@@ -73,17 +73,8 @@ synchronized void parseHeartbeat(final int nodeId, final PipeHeartbeat pipeHeart
7373
if (heartbeatCount % registeredNodeNumber == 0) {
7474
canSubmitHandleMetaChangeProcedure.set(true);
7575

76-
// registeredNodeNumber may be changed, update it here when we can submit procedure
77-
registeredNodeNumber = configManager.getNodeManager().getRegisteredNodeCount();
78-
if (registeredNodeNumber <= 0) {
79-
LOGGER.warn(
80-
"registeredNodeNumber is {} when parseHeartbeat from node (id={}).",
81-
registeredNodeNumber,
82-
nodeId);
83-
// registeredNodeNumber can not be set to 0 in this class, otherwise may cause
84-
// DivideByZeroException
85-
registeredNodeNumber = 1;
86-
}
76+
// The expected reporter set may be changed, update it at the end of the current round.
77+
registeredNodeNumber = getExpectedHeartbeatNodeCount();
8778
}
8879

8980
if (pipeHeartbeat.isEmpty()
@@ -114,21 +105,32 @@ synchronized void parseHeartbeat(final int nodeId, final PipeHeartbeat pipeHeart
114105
if (canSubmitHandleMetaChangeProcedure.get()
115106
&& (needWriteConsensusOnConfigNodes.get()
116107
|| needPushPipeMetaToDataNodes.get())) {
117-
configManager
108+
if (configManager
118109
.getProcedureManager()
119110
.pipeHandleMetaChange(
120-
needWriteConsensusOnConfigNodes.get(), needPushPipeMetaToDataNodes.get());
121-
122-
// Reset flags after procedure is submitted
123-
needWriteConsensusOnConfigNodes.set(false);
124-
needPushPipeMetaToDataNodes.set(false);
111+
needWriteConsensusOnConfigNodes.get(),
112+
needPushPipeMetaToDataNodes.get())) {
113+
needWriteConsensusOnConfigNodes.set(false);
114+
needPushPipeMetaToDataNodes.set(false);
115+
}
125116
}
126117
} finally {
127118
configManager.getPipeManager().getPipeTaskCoordinator().unlock();
128119
}
129120
});
130121
}
131122

123+
private int getExpectedHeartbeatNodeCount() {
124+
final int expectedNodeCount =
125+
configManager.getNodeManager().getRegisteredDataNodeCount()
126+
+ (PipeConfig.getInstance().isSeperatedPipeHeartbeatEnabled() ? 1 : 0);
127+
if (expectedNodeCount <= 0) {
128+
LOGGER.warn("Expected pipe heartbeat node count is {}, fallback to 1.", expectedNodeCount);
129+
return 1;
130+
}
131+
return expectedNodeCount;
132+
}
133+
132134
private void parseHeartbeatAndSaveMetaChangeLocally(
133135
final AtomicReference<PipeTaskInfo> pipeTaskInfo,
134136
final int nodeId,

iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/runtime/PipeHandleLeaderChangeProcedure.java

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -164,7 +164,11 @@ public void deserialize(ByteBuffer byteBuffer) {
164164
final int oldDataRegionLeaderId = ReadWriteIOUtils.readInt(byteBuffer);
165165
final int newDataRegionLeaderId = ReadWriteIOUtils.readInt(byteBuffer);
166166
regionGroupToOldAndNewLeaderPairMap.put(
167-
new TConsensusGroupId(TConsensusGroupType.DataRegion, dataRegionGroupId),
167+
new TConsensusGroupId(
168+
dataRegionGroupId == Integer.MIN_VALUE
169+
? TConsensusGroupType.ConfigRegion
170+
: TConsensusGroupType.DataRegion,
171+
dataRegionGroupId),
168172
new Pair<>(oldDataRegionLeaderId, newDataRegionLeaderId));
169173
}
170174
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,182 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
20+
package org.apache.iotdb.confignode.manager.pipe.coordinator.runtime.heartbeat;
21+
22+
import org.apache.iotdb.commons.conf.CommonDescriptor;
23+
import org.apache.iotdb.confignode.manager.ConfigManager;
24+
import org.apache.iotdb.confignode.manager.ProcedureManager;
25+
import org.apache.iotdb.confignode.manager.node.NodeManager;
26+
import org.apache.iotdb.confignode.manager.pipe.coordinator.PipeManager;
27+
import org.apache.iotdb.confignode.manager.pipe.coordinator.runtime.PipeRuntimeCoordinator;
28+
import org.apache.iotdb.confignode.manager.pipe.coordinator.task.PipeTaskCoordinator;
29+
import org.apache.iotdb.confignode.persistence.pipe.PipeTaskInfo;
30+
31+
import org.junit.After;
32+
import org.junit.Before;
33+
import org.junit.Test;
34+
import org.mockito.Mockito;
35+
36+
import java.lang.reflect.Field;
37+
import java.util.Collections;
38+
import java.util.concurrent.CompletableFuture;
39+
import java.util.concurrent.ExecutorService;
40+
import java.util.concurrent.atomic.AtomicBoolean;
41+
import java.util.concurrent.atomic.AtomicReference;
42+
43+
import static org.mockito.ArgumentMatchers.any;
44+
import static org.mockito.ArgumentMatchers.anyBoolean;
45+
import static org.mockito.Mockito.never;
46+
import static org.mockito.Mockito.times;
47+
import static org.mockito.Mockito.verify;
48+
import static org.mockito.Mockito.when;
49+
50+
public class PipeHeartbeatParserTest {
51+
52+
private boolean originalSeparatedPipeHeartbeatEnabled;
53+
54+
@Before
55+
public void setUp() {
56+
originalSeparatedPipeHeartbeatEnabled =
57+
CommonDescriptor.getInstance().getConfig().isSeperatedPipeHeartbeatEnabled();
58+
}
59+
60+
@After
61+
public void tearDown() {
62+
CommonDescriptor.getInstance()
63+
.getConfig()
64+
.setSeperatedPipeHeartbeatEnabled(originalSeparatedPipeHeartbeatEnabled);
65+
}
66+
67+
@Test
68+
public void testParseHeartbeatCountsOnlyDataNodesWhenSeparatedHeartbeatDisabled()
69+
throws Exception {
70+
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
71+
72+
final ParserTestContext context = createParserTestContext(2);
73+
setMetaChangeFlags(context.parser, true, false);
74+
75+
context.parser.parseHeartbeat(1, emptyHeartbeat());
76+
verify(context.procedureManager, never()).pipeHandleMetaChange(anyBoolean(), anyBoolean());
77+
78+
context.parser.parseHeartbeat(2, emptyHeartbeat());
79+
verify(context.procedureManager, times(1)).pipeHandleMetaChange(true, false);
80+
}
81+
82+
@Test
83+
public void testParseHeartbeatCountsLocalConfigNodeWhenSeparatedHeartbeatEnabled()
84+
throws Exception {
85+
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(true);
86+
87+
final ParserTestContext context = createParserTestContext(2);
88+
setMetaChangeFlags(context.parser, true, false);
89+
90+
context.parser.parseHeartbeat(1, emptyHeartbeat());
91+
context.parser.parseHeartbeat(2, emptyHeartbeat());
92+
verify(context.procedureManager, never()).pipeHandleMetaChange(anyBoolean(), anyBoolean());
93+
94+
context.parser.parseHeartbeat(3, emptyHeartbeat());
95+
verify(context.procedureManager, times(1)).pipeHandleMetaChange(true, false);
96+
}
97+
98+
@Test
99+
public void testParseHeartbeatKeepsPendingFlagsWhenProcedureSubmissionFails() throws Exception {
100+
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
101+
102+
final ParserTestContext context = createParserTestContext(2);
103+
when(context.procedureManager.pipeHandleMetaChange(anyBoolean(), anyBoolean()))
104+
.thenReturn(false, true);
105+
setMetaChangeFlags(context.parser, true, false);
106+
107+
context.parser.parseHeartbeat(1, emptyHeartbeat());
108+
verify(context.procedureManager, never()).pipeHandleMetaChange(anyBoolean(), anyBoolean());
109+
110+
context.parser.parseHeartbeat(2, emptyHeartbeat());
111+
verify(context.procedureManager, times(1)).pipeHandleMetaChange(true, false);
112+
113+
context.parser.parseHeartbeat(3, emptyHeartbeat());
114+
verify(context.procedureManager, times(1)).pipeHandleMetaChange(true, false);
115+
116+
context.parser.parseHeartbeat(4, emptyHeartbeat());
117+
verify(context.procedureManager, times(2)).pipeHandleMetaChange(true, false);
118+
}
119+
120+
private ParserTestContext createParserTestContext(final int registeredDataNodeCount) {
121+
final ConfigManager configManager = Mockito.mock(ConfigManager.class);
122+
final NodeManager nodeManager = Mockito.mock(NodeManager.class);
123+
final ProcedureManager procedureManager = Mockito.mock(ProcedureManager.class);
124+
final PipeManager pipeManager = Mockito.mock(PipeManager.class);
125+
final PipeRuntimeCoordinator pipeRuntimeCoordinator =
126+
Mockito.mock(PipeRuntimeCoordinator.class);
127+
final PipeTaskCoordinator pipeTaskCoordinator = Mockito.mock(PipeTaskCoordinator.class);
128+
final ExecutorService procedureSubmitter = Mockito.mock(ExecutorService.class);
129+
130+
when(configManager.getNodeManager()).thenReturn(nodeManager);
131+
when(configManager.getProcedureManager()).thenReturn(procedureManager);
132+
when(configManager.getPipeManager()).thenReturn(pipeManager);
133+
when(nodeManager.getRegisteredDataNodeCount()).thenReturn(registeredDataNodeCount);
134+
when(pipeManager.getPipeRuntimeCoordinator()).thenReturn(pipeRuntimeCoordinator);
135+
when(pipeManager.getPipeTaskCoordinator()).thenReturn(pipeTaskCoordinator);
136+
when(pipeRuntimeCoordinator.getProcedureSubmitter()).thenReturn(procedureSubmitter);
137+
when(pipeTaskCoordinator.tryLock()).thenReturn(new AtomicReference<>(new PipeTaskInfo()));
138+
when(procedureManager.pipeHandleMetaChange(anyBoolean(), anyBoolean())).thenReturn(true);
139+
Mockito.doAnswer(
140+
invocation -> {
141+
((Runnable) invocation.getArgument(0)).run();
142+
return CompletableFuture.completedFuture(null);
143+
})
144+
.when(procedureSubmitter)
145+
.submit(any(Runnable.class));
146+
147+
return new ParserTestContext(new PipeHeartbeatParser(configManager), procedureManager);
148+
}
149+
150+
private void setMetaChangeFlags(
151+
final PipeHeartbeatParser parser,
152+
final boolean needWriteConsensusOnConfigNodes,
153+
final boolean needPushPipeMetaToDataNodes)
154+
throws Exception {
155+
setAtomicBooleanField(
156+
parser, "needWriteConsensusOnConfigNodes", needWriteConsensusOnConfigNodes);
157+
setAtomicBooleanField(parser, "needPushPipeMetaToDataNodes", needPushPipeMetaToDataNodes);
158+
}
159+
160+
private void setAtomicBooleanField(
161+
final PipeHeartbeatParser parser, final String fieldName, final boolean value)
162+
throws Exception {
163+
final Field field = PipeHeartbeatParser.class.getDeclaredField(fieldName);
164+
field.setAccessible(true);
165+
((AtomicBoolean) field.get(parser)).set(value);
166+
}
167+
168+
private PipeHeartbeat emptyHeartbeat() {
169+
return new PipeHeartbeat(Collections.emptyList(), null, null, null);
170+
}
171+
172+
private static class ParserTestContext {
173+
private final PipeHeartbeatParser parser;
174+
private final ProcedureManager procedureManager;
175+
176+
private ParserTestContext(
177+
final PipeHeartbeatParser parser, final ProcedureManager procedureManager) {
178+
this.parser = parser;
179+
this.procedureManager = procedureManager;
180+
}
181+
}
182+
}

iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/pipe/runtime/PipeHandleLeaderChangeProcedureTest.java

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,10 @@
2121

2222
import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
2323
import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
24+
import org.apache.iotdb.confignode.procedure.Procedure;
25+
import org.apache.iotdb.confignode.procedure.state.ProcedureState;
2426
import org.apache.iotdb.confignode.procedure.store.ProcedureFactory;
27+
import org.apache.iotdb.confignode.procedure.store.ProcedureType;
2528

2629
import org.apache.tsfile.utils.Pair;
2730
import org.apache.tsfile.utils.PublicBAOS;
@@ -45,6 +48,9 @@ public void serializeDeserializeTest() {
4548
leaderMap.put(new TConsensusGroupId(TConsensusGroupType.DataRegion, 1), new Pair<>(1, 2));
4649
leaderMap.put(new TConsensusGroupId(TConsensusGroupType.DataRegion, 2), new Pair<>(2, 3));
4750
leaderMap.put(new TConsensusGroupId(TConsensusGroupType.DataRegion, 3), new Pair<>(4, 5));
51+
leaderMap.put(
52+
new TConsensusGroupId(TConsensusGroupType.ConfigRegion, Integer.MIN_VALUE),
53+
new Pair<>(6, 7));
4854

4955
PipeHandleLeaderChangeProcedure proc = new PipeHandleLeaderChangeProcedure(leaderMap);
5056

@@ -60,4 +66,44 @@ public void serializeDeserializeTest() {
6066
fail();
6167
}
6268
}
69+
70+
@Test
71+
public void deserializeOldFormatConfigRegionTest() {
72+
PublicBAOS byteArrayOutputStream = new PublicBAOS();
73+
DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream);
74+
75+
Map<TConsensusGroupId, Pair<Integer, Integer>> leaderMap = new HashMap<>();
76+
leaderMap.put(
77+
new TConsensusGroupId(TConsensusGroupType.ConfigRegion, Integer.MIN_VALUE),
78+
new Pair<>(6, 7));
79+
80+
try {
81+
outputStream.writeShort(ProcedureType.PIPE_HANDLE_LEADER_CHANGE_PROCEDURE.getTypeCode());
82+
outputStream.writeLong(Procedure.NO_PROC_ID);
83+
outputStream.writeInt(ProcedureState.INITIALIZING.ordinal());
84+
outputStream.writeLong(0L);
85+
outputStream.writeLong(0L);
86+
outputStream.writeLong(Procedure.NO_PROC_ID);
87+
outputStream.writeLong(Procedure.NO_TIMEOUT);
88+
outputStream.writeInt(-1);
89+
outputStream.write((byte) 0);
90+
outputStream.writeInt(-1);
91+
outputStream.write((byte) 0);
92+
outputStream.writeInt(0);
93+
outputStream.write((byte) 0);
94+
outputStream.writeInt(leaderMap.size());
95+
outputStream.writeInt(Integer.MIN_VALUE);
96+
outputStream.writeInt(6);
97+
outputStream.writeInt(7);
98+
99+
ByteBuffer buffer =
100+
ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, byteArrayOutputStream.size());
101+
PipeHandleLeaderChangeProcedure proc =
102+
(PipeHandleLeaderChangeProcedure) ProcedureFactory.getInstance().create(buffer);
103+
104+
assertEquals(new PipeHandleLeaderChangeProcedure(leaderMap), proc);
105+
} catch (Exception e) {
106+
fail();
107+
}
108+
}
63109
}

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/airgap/IoTDBAirGapReceiver.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -178,6 +178,11 @@ private void handleReq(final AirGapPseudoTPipeTransferRequest req, final long st
178178
if (System.currentTimeMillis() - startTime
179179
< PipeConfig.getInstance().getPipeAirGapRetryMaxMs()) {
180180
handleReq(req, startTime);
181+
} else {
182+
LOGGER.warn(
183+
"Pipe air gap receiver {}: Temporary unavailable retry timed out, returning FAIL to sender.",
184+
receiverId);
185+
fail();
181186
}
182187
} else {
183188
LOGGER.warn(

0 commit comments

Comments
 (0)