Skip to content

Commit ca52520

Browse files
authored
[To dev/1.3] Fix pipe drop event discard with restart-aware committer keys (#17748) (#17778)
1 parent 0f253d8 commit ca52520

18 files changed

Lines changed: 171 additions & 96 deletions

File tree

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -202,7 +202,10 @@ private void collectEvent(final Event event) {
202202
enrichedEvent.setRebootTimes(PipeDataNodeAgent.runtime().getRebootTimes());
203203

204204
if (enrichedEvent.getPipeName() != null
205-
&& pendingQueue.isPipeDropped(enrichedEvent.getPipeName(), creationTime, regionId)) {
205+
&& (pendingQueue.isEventFromDroppedPipe(enrichedEvent)
206+
|| (enrichedEvent.getCommitterKey() == null
207+
&& pendingQueue.isPipeDropped(
208+
enrichedEvent.getPipeName(), creationTime, regionId)))) {
206209
enrichedEvent.clearReferenceCount(PipeEventCollector.class.getName());
207210
return;
208211
}

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeRealtimePriorityBlockingQueue.java

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -360,12 +360,16 @@ public void discardAllEvents() {
360360
@Override
361361
public void discardEventsOfPipe(
362362
final String pipeNameToDrop, final long creationTimeToDrop, final int regionId) {
363-
super.discardEventsOfPipe(pipeNameToDrop, creationTimeToDrop, regionId);
363+
discardEventsOfPipe(new CommitterKey(pipeNameToDrop, creationTimeToDrop, regionId, -1));
364+
}
365+
366+
@Override
367+
public void discardEventsOfPipe(final CommitterKey committerKey) {
368+
super.discardEventsOfPipe(committerKey);
364369
tsfileInsertEventDeque.removeIf(
365370
event -> {
366371
if (event instanceof EnrichedEvent
367-
&& isEventFromPipe(
368-
((EnrichedEvent) event), pipeNameToDrop, creationTimeToDrop, regionId)) {
372+
&& isEventFromPipe((EnrichedEvent) event, committerKey)) {
369373
if (((EnrichedEvent) event)
370374
.clearReferenceCount(PipeRealtimePriorityBlockingQueue.class.getName())) {
371375
eventCounter.decreaseEventCount(event);

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtask.java

Lines changed: 14 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121

2222
import org.apache.iotdb.commons.exception.pipe.PipeRuntimeException;
2323
import org.apache.iotdb.commons.pipe.agent.task.connection.UnboundedBlockingPendingQueue;
24+
import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey;
2425
import org.apache.iotdb.commons.pipe.agent.task.subtask.PipeAbstractSinkSubtask;
2526
import org.apache.iotdb.commons.pipe.config.PipeConfig;
2627
import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
@@ -201,10 +202,9 @@ public void close() {
201202
* When a pipe is dropped, the connector maybe reused and will not be closed. So we just discard
202203
* its queued events in the output pipe connector.
203204
*/
204-
public void discardEventsOfPipe(
205-
final String pipeNameToDrop, final long creationTimeToDrop, final int regionId) {
205+
public void discardEventsOfPipe(final CommitterKey committerKey) {
206206
// Try to remove the events as much as possible
207-
inputPendingQueue.discardEventsOfPipe(pipeNameToDrop, creationTimeToDrop, regionId);
207+
inputPendingQueue.discardEventsOfPipe(committerKey);
208208

209209
try {
210210
increaseHighPriorityTaskCount();
@@ -217,9 +217,7 @@ public void discardEventsOfPipe(
217217
// use a new thread to stop all the pipes, we will not encounter deadlock here. Or else we
218218
// will.
219219
if (lastEvent instanceof EnrichedEvent
220-
&& pipeNameToDrop.equals(((EnrichedEvent) lastEvent).getPipeName())
221-
&& creationTimeToDrop == ((EnrichedEvent) lastEvent).getCreationTime()
222-
&& regionId == ((EnrichedEvent) lastEvent).getRegionId()) {
220+
&& isEventFromPipe((EnrichedEvent) lastEvent, committerKey)) {
223221
// Do not clear the last event's reference counts because it may be on transferring
224222
lastEvent = null;
225223
// Submit self to avoid that the lastEvent has been retried "max times" times and has
@@ -241,9 +239,7 @@ public void discardEventsOfPipe(
241239
// clear the lastExceptionEvent. It's safe to potentially clear it twice because we have the
242240
// "nonnull" detection.
243241
if (lastExceptionEvent instanceof EnrichedEvent
244-
&& pipeNameToDrop.equals(((EnrichedEvent) lastExceptionEvent).getPipeName())
245-
&& creationTimeToDrop == ((EnrichedEvent) lastExceptionEvent).getCreationTime()
246-
&& regionId == ((EnrichedEvent) lastExceptionEvent).getRegionId()) {
242+
&& isEventFromPipe((EnrichedEvent) lastExceptionEvent, committerKey)) {
247243
clearReferenceCountAndReleaseLastExceptionEvent();
248244
}
249245
}
@@ -252,11 +248,18 @@ public void discardEventsOfPipe(
252248
}
253249

254250
if (outputPipeConnector instanceof PipeConnectorWithEventDiscard) {
255-
((PipeConnectorWithEventDiscard) outputPipeConnector)
256-
.discardEventsOfPipe(pipeNameToDrop, creationTimeToDrop, regionId);
251+
((PipeConnectorWithEventDiscard) outputPipeConnector).discardEventsOfPipe(committerKey);
257252
}
258253
}
259254

255+
private static boolean isEventFromPipe(
256+
final EnrichedEvent event, final CommitterKey committerKey) {
257+
return committerKey.getPipeName().equals(event.getPipeName())
258+
&& committerKey.getCreationTime() == event.getCreationTime()
259+
&& committerKey.getRegionId() == event.getRegionId()
260+
&& (committerKey.getRestartTimes() < 0 || committerKey.equals(event.getCommitterKey()));
261+
}
262+
260263
//////////////////////////// APIs provided for metric framework ////////////////////////////
261264

262265
public String getAttributeSortedString() {

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskLifeCycle.java

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
package org.apache.iotdb.db.pipe.agent.task.subtask.sink;
2121

2222
import org.apache.iotdb.commons.pipe.agent.task.connection.UnboundedBlockingPendingQueue;
23+
import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey;
2324
import org.apache.iotdb.db.pipe.agent.task.execution.PipeSinkSubtaskExecutor;
2425
import org.apache.iotdb.db.pipe.resource.PipeDataNodeResourceManager;
2526
import org.apache.iotdb.pipe.api.event.Event;
@@ -86,19 +87,17 @@ public synchronized void register() {
8687
* Otherwise, the {@link PipeSinkSubtaskLifeCycle#runningTaskCount} might be inconsistent with the
8788
* {@link PipeSinkSubtaskLifeCycle#registeredTaskCount} because of parallel connector scheduling.
8889
*
89-
* @param pipeNameToDeregister pipe name
90-
* @param regionId region id
90+
* @param committerKey committer key of the pipe task to deregister
9191
* @return {@code true} if the {@link PipeSinkSubtask} is out of life cycle, indicating that the
9292
* {@link PipeSinkSubtask} should never be used again
9393
* @throws IllegalStateException if {@link PipeSinkSubtaskLifeCycle#registeredTaskCount} <= 0
9494
*/
95-
public synchronized boolean deregister(
96-
final String pipeNameToDeregister, final long creationTimeToDeregister, final int regionId) {
95+
public synchronized boolean deregister(final CommitterKey committerKey) {
9796
if (registeredTaskCount <= 0) {
9897
throw new IllegalStateException("registeredTaskCount <= 0");
9998
}
10099

101-
subtask.discardEventsOfPipe(pipeNameToDeregister, creationTimeToDeregister, regionId);
100+
subtask.discardEventsOfPipe(committerKey);
102101

103102
try {
104103
if (registeredTaskCount > 1) {

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskManager.java

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
import org.apache.iotdb.commons.consensus.DataRegionId;
2323
import org.apache.iotdb.commons.pipe.agent.plugin.builtin.BuiltinPipePlugin;
2424
import org.apache.iotdb.commons.pipe.agent.task.connection.UnboundedBlockingPendingQueue;
25+
import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey;
2526
import org.apache.iotdb.commons.pipe.agent.task.progress.PipeEventCommitManager;
2627
import org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant;
2728
import org.apache.iotdb.commons.pipe.config.constant.SystemConstant;
@@ -209,7 +210,10 @@ public synchronized void deregister(
209210
// Shall not be empty
210211
final PipeSinkSubtaskExecutor executor = lifeCycles.get(0).executor;
211212

212-
lifeCycles.removeIf(o -> o.deregister(pipeName, creationTime, regionId));
213+
final CommitterKey committerKey =
214+
PipeEventCommitManager.getInstance().getCommitterKey(pipeName, creationTime, regionId);
215+
216+
lifeCycles.removeIf(o -> o.deregister(committerKey));
213217

214218
if (lifeCycles.isEmpty()) {
215219
attributeSortedString2SubtaskLifeCycleMap.remove(attributeSortedString);

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/batch/PipeTabletEventBatch.java

Lines changed: 14 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919

2020
package org.apache.iotdb.db.pipe.sink.payload.evolvable.batch;
2121

22+
import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey;
2223
import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
2324
import org.apache.iotdb.db.pipe.resource.PipeDataNodeResourceManager;
2425
import org.apache.iotdb.db.pipe.resource.memory.PipeMemoryBlock;
@@ -156,18 +157,28 @@ public synchronized void close() {
156157
*/
157158
public synchronized void discardEventsOfPipe(
158159
final String pipeNameToDrop, final long creationTimeToDrop, final int regionId) {
160+
discardEventsOfPipe(new CommitterKey(pipeNameToDrop, creationTimeToDrop, regionId, -1));
161+
}
162+
163+
public synchronized void discardEventsOfPipe(final CommitterKey committerKey) {
159164
events.removeIf(
160165
event -> {
161-
if (pipeNameToDrop.equals(event.getPipeName())
162-
&& creationTimeToDrop == event.getCreationTime()
163-
&& regionId == event.getRegionId()) {
166+
if (isEventFromPipe(event, committerKey)) {
164167
event.clearReferenceCount(IoTDBDataRegionAsyncSink.class.getName());
165168
return true;
166169
}
167170
return false;
168171
});
169172
}
170173

174+
private static boolean isEventFromPipe(
175+
final EnrichedEvent event, final CommitterKey committerKey) {
176+
return committerKey.getPipeName().equals(event.getPipeName())
177+
&& committerKey.getCreationTime() == event.getCreationTime()
178+
&& committerKey.getRegionId() == event.getRegionId()
179+
&& (committerKey.getRestartTimes() < 0 || committerKey.equals(event.getCommitterKey()));
180+
}
181+
171182
public synchronized void decreaseEventsReferenceCount(
172183
final String holderMessage, final boolean shouldReport) {
173184
events.forEach(event -> event.decreaseReferenceCount(holderMessage, shouldReport));

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/batch/PipeTransferBatchReqBuilder.java

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
package org.apache.iotdb.db.pipe.sink.payload.evolvable.batch;
2121

2222
import org.apache.iotdb.common.rpc.thrift.TEndPoint;
23+
import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey;
2324
import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
2425
import org.apache.iotdb.db.pipe.event.common.tablet.PipeInsertNodeTabletInsertionEvent;
2526
import org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent;
@@ -197,10 +198,12 @@ public boolean isEmpty() {
197198

198199
public synchronized void discardEventsOfPipe(
199200
final String pipeNameToDrop, final long creationTimeToDrop, final int regionId) {
200-
defaultBatch.discardEventsOfPipe(pipeNameToDrop, creationTimeToDrop, regionId);
201-
endPointToBatch
202-
.values()
203-
.forEach(batch -> batch.discardEventsOfPipe(pipeNameToDrop, creationTimeToDrop, regionId));
201+
discardEventsOfPipe(new CommitterKey(pipeNameToDrop, creationTimeToDrop, regionId, -1));
202+
}
203+
204+
public synchronized void discardEventsOfPipe(final CommitterKey committerKey) {
205+
defaultBatch.discardEventsOfPipe(committerKey);
206+
endPointToBatch.values().forEach(batch -> batch.discardEventsOfPipe(committerKey));
204207
}
205208

206209
public int size() {

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/airgap/IoTDBDataRegionAirGapSink.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
package org.apache.iotdb.db.pipe.sink.protocol.airgap;
2121

2222
import org.apache.iotdb.common.rpc.thrift.TSStatus;
23+
import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey;
2324
import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
2425
import org.apache.iotdb.commons.pipe.sink.limiter.TsFileSendRateLimiter;
2526
import org.apache.iotdb.commons.utils.RetryUtils;
@@ -546,8 +547,13 @@ protected byte[] compressIfNeeded(final byte[] reqInBytes) throws IOException {
546547
@Override
547548
public synchronized void discardEventsOfPipe(
548549
final String pipeNameToDrop, final long creationTimeToDrop, final int regionId) {
550+
discardEventsOfPipe(new CommitterKey(pipeNameToDrop, creationTimeToDrop, regionId, -1));
551+
}
552+
553+
@Override
554+
public synchronized void discardEventsOfPipe(final CommitterKey committerKey) {
549555
if (Objects.nonNull(tabletBatchBuilder)) {
550-
tabletBatchBuilder.discardEventsOfPipe(pipeNameToDrop, creationTimeToDrop, regionId);
556+
tabletBatchBuilder.discardEventsOfPipe(committerKey);
551557
}
552558
}
553559

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSink.java

Lines changed: 19 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -22,8 +22,8 @@
2222
import org.apache.iotdb.common.rpc.thrift.TEndPoint;
2323
import org.apache.iotdb.commons.client.ThriftClient;
2424
import org.apache.iotdb.commons.client.async.AsyncPipeDataTransferServiceClient;
25+
import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey;
2526
import org.apache.iotdb.commons.pipe.config.PipeConfig;
26-
import org.apache.iotdb.commons.pipe.datastructure.Triple;
2727
import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
2828
import org.apache.iotdb.commons.pipe.resource.log.PipeLogger;
2929
import org.apache.iotdb.commons.pipe.sink.protocol.IoTDBSink;
@@ -123,9 +123,7 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
123123
private final Map<PipeTransferTrackableHandler, PipeTransferTrackableHandler> pendingHandlers =
124124
new ConcurrentHashMap<>();
125125

126-
// Pipe name, creation time, region id
127-
private final Set<Triple<String, Long, Integer>> droppedPipeTaskKeys =
128-
ConcurrentHashMap.newKeySet();
126+
private final Set<CommitterKey> droppedPipeTaskKeys = ConcurrentHashMap.newKeySet();
129127

130128
private boolean enableSendTsFileLimit;
131129
private volatile boolean isConnectionException;
@@ -722,16 +720,20 @@ public boolean isEnableSendTsFileLimit() {
722720
@Override
723721
public synchronized void discardEventsOfPipe(
724722
final String pipeNameToDrop, final long creationTimeToDrop, final int regionId) {
725-
droppedPipeTaskKeys.add(new Triple<>(pipeNameToDrop, creationTimeToDrop, regionId));
723+
discardEventsOfPipe(new CommitterKey(pipeNameToDrop, creationTimeToDrop, regionId, -1));
724+
}
725+
726+
@Override
727+
public synchronized void discardEventsOfPipe(final CommitterKey committerKey) {
728+
droppedPipeTaskKeys.add(committerKey);
726729

727730
if (isTabletBatchModeEnabled && Objects.nonNull(tabletBatchBuilder)) {
728-
tabletBatchBuilder.discardEventsOfPipe(pipeNameToDrop, creationTimeToDrop, regionId);
731+
tabletBatchBuilder.discardEventsOfPipe(committerKey);
729732
}
730733
retryEventQueue.removeIf(
731734
event -> {
732735
if (event instanceof EnrichedEvent
733-
&& isDroppedPipe(
734-
(EnrichedEvent) event, pipeNameToDrop, creationTimeToDrop, regionId)) {
736+
&& isDroppedPipe((EnrichedEvent) event, committerKey)) {
735737
((EnrichedEvent) event).clearReferenceCount(IoTDBDataRegionAsyncSink.class.getName());
736738
retryEventQueueEventCounter.decreaseEventCount(event);
737739
return true;
@@ -742,8 +744,7 @@ && isDroppedPipe(
742744
retryTsFileQueue.removeIf(
743745
event -> {
744746
if (event instanceof EnrichedEvent
745-
&& isDroppedPipe(
746-
(EnrichedEvent) event, pipeNameToDrop, creationTimeToDrop, regionId)) {
747+
&& isDroppedPipe((EnrichedEvent) event, committerKey)) {
747748
((EnrichedEvent) event).clearReferenceCount(IoTDBDataRegionAsyncSink.class.getName());
748749
retryEventQueueEventCounter.decreaseEventCount(event);
749750
return true;
@@ -845,18 +846,14 @@ public void setTransferTsFileCounter(AtomicInteger transferTsFileCounter) {
845846
}
846847

847848
private boolean isDroppedPipe(final EnrichedEvent event) {
848-
return droppedPipeTaskKeys.contains(
849-
new Triple<>(event.getPipeName(), event.getCreationTime(), event.getRegionId()));
850-
}
851-
852-
private static boolean isDroppedPipe(
853-
final EnrichedEvent event,
854-
final String pipeNameToDrop,
855-
final long creationTimeToDrop,
856-
final int regionId) {
857-
return pipeNameToDrop.equals(event.getPipeName())
858-
&& creationTimeToDrop == event.getCreationTime()
859-
&& regionId == event.getRegionId();
849+
return droppedPipeTaskKeys.stream().anyMatch(key -> isDroppedPipe(event, key));
850+
}
851+
852+
private static boolean isDroppedPipe(final EnrichedEvent event, final CommitterKey committerKey) {
853+
return committerKey.getPipeName().equals(event.getPipeName())
854+
&& committerKey.getCreationTime() == event.getCreationTime()
855+
&& committerKey.getRegionId() == event.getRegionId()
856+
&& (committerKey.getRestartTimes() < 0 || committerKey.equals(event.getCommitterKey()));
860857
}
861858

862859
@Override

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/sync/IoTDBDataRegionSyncSink.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121

2222
import org.apache.iotdb.common.rpc.thrift.TEndPoint;
2323
import org.apache.iotdb.common.rpc.thrift.TSStatus;
24+
import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey;
2425
import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
2526
import org.apache.iotdb.commons.pipe.sink.client.IoTDBSyncClient;
2627
import org.apache.iotdb.commons.pipe.sink.limiter.TsFileSendRateLimiter;
@@ -523,8 +524,13 @@ public TPipeTransferReq compressIfNeeded(final TPipeTransferReq req) throws IOEx
523524
@Override
524525
public synchronized void discardEventsOfPipe(
525526
final String pipeNameToDrop, final long creationTimeToDrop, final int regionId) {
527+
discardEventsOfPipe(new CommitterKey(pipeNameToDrop, creationTimeToDrop, regionId, -1));
528+
}
529+
530+
@Override
531+
public synchronized void discardEventsOfPipe(final CommitterKey committerKey) {
526532
if (Objects.nonNull(tabletBatchBuilder)) {
527-
tabletBatchBuilder.discardEventsOfPipe(pipeNameToDrop, creationTimeToDrop, regionId);
533+
tabletBatchBuilder.discardEventsOfPipe(committerKey);
528534
}
529535
}
530536

0 commit comments

Comments
 (0)