Skip to content

Commit afb9059

Browse files
authored
Fix subscription TsFile deduplication and poll concurrency (#18758)
1 parent bd3109e commit afb9059

3 files changed

Lines changed: 76 additions & 6 deletions

File tree

‎iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -135,7 +135,8 @@ abstract class AbstractSubscriptionConsumer implements AutoCloseable {
135135

136136
private final String fileSaveDir;
137137
private final boolean fileSaveFsync;
138-
private final Set<SubscriptionCommitContext> inFlightFilesCommitContextSet = new HashSet<>();
138+
private final Set<SubscriptionCommitContext> inFlightFilesCommitContextSet =
139+
ConcurrentHashMap.newKeySet();
139140

140141
private final int thriftMaxFrameSize;
141142
private final int connectionTimeoutInMs;

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/TsFileDeduplicationBlockingPendingQueue.java‎

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -40,13 +40,13 @@ public class TsFileDeduplicationBlockingPendingQueue extends SubscriptionBlockin
4040
private static final Logger LOGGER =
4141
LoggerFactory.getLogger(TsFileDeduplicationBlockingPendingQueue.class);
4242

43-
private final Cache<Integer, Boolean> hashCodeToIsGeneratedByHistoricalExtractor;
43+
private final Cache<String, Boolean> tsFilePathToIsGeneratedByHistoricalExtractor;
4444

4545
public TsFileDeduplicationBlockingPendingQueue(
4646
final UnboundedBlockingPendingQueue<Event> inputPendingQueue) {
4747
super(inputPendingQueue);
4848

49-
this.hashCodeToIsGeneratedByHistoricalExtractor =
49+
this.tsFilePathToIsGeneratedByHistoricalExtractor =
5050
Caffeine.newBuilder()
5151
.expireAfterAccess(
5252
SubscriptionConfig.getInstance().getSubscriptionTsFileDeduplicationWindowSeconds(),
@@ -101,12 +101,13 @@ && isDuplicated((PipeTsFileInsertionEvent) sourceEvent)) {
101101
}
102102

103103
private boolean isDuplicated(final PipeTsFileInsertionEvent event) {
104-
final int hashCode = event.getTsFile().hashCode();
104+
final String tsFilePath = event.getTsFile().toPath().toAbsolutePath().normalize().toString();
105105
final boolean isGeneratedByHistoricalExtractor = event.isGeneratedByHistoricalExtractor();
106106
final Boolean existedIsGeneratedByHistoricalExtractor =
107-
hashCodeToIsGeneratedByHistoricalExtractor.getIfPresent(hashCode);
107+
tsFilePathToIsGeneratedByHistoricalExtractor.getIfPresent(tsFilePath);
108108
if (Objects.isNull(existedIsGeneratedByHistoricalExtractor)) {
109-
hashCodeToIsGeneratedByHistoricalExtractor.put(hashCode, isGeneratedByHistoricalExtractor);
109+
tsFilePathToIsGeneratedByHistoricalExtractor.put(
110+
tsFilePath, isGeneratedByHistoricalExtractor);
110111
return false;
111112
}
112113
// Multiple PipeRawTabletInsertionEvents parsed from the same PipeTsFileInsertionEvent (i.e.,
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,68 @@
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.db.subscription.broker;
21+
22+
import org.apache.iotdb.commons.pipe.agent.task.connection.UnboundedBlockingPendingQueue;
23+
import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
24+
import org.apache.iotdb.db.pipe.metric.source.PipeDataRegionEventCounter;
25+
import org.apache.iotdb.pipe.api.event.Event;
26+
27+
import org.junit.Assert;
28+
import org.junit.Test;
29+
import org.mockito.Mockito;
30+
31+
import java.io.File;
32+
33+
public class TsFileDeduplicationBlockingPendingQueueTest {
34+
35+
@Test
36+
public void testDifferentPathsWithSameHashCodeAreNotDeduplicated() {
37+
final UnboundedBlockingPendingQueue<Event> inputPendingQueue =
38+
new UnboundedBlockingPendingQueue<>(new PipeDataRegionEventCounter());
39+
final TsFileDeduplicationBlockingPendingQueue queue =
40+
new TsFileDeduplicationBlockingPendingQueue(inputPendingQueue);
41+
42+
final PipeTsFileInsertionEvent firstEvent = Mockito.mock(PipeTsFileInsertionEvent.class);
43+
Mockito.when(firstEvent.getTsFile()).thenReturn(new CollidingFile("first.tsfile"));
44+
Mockito.when(firstEvent.isGeneratedByHistoricalExtractor()).thenReturn(false);
45+
46+
final PipeTsFileInsertionEvent secondEvent = Mockito.mock(PipeTsFileInsertionEvent.class);
47+
Mockito.when(secondEvent.getTsFile()).thenReturn(new CollidingFile("second.tsfile"));
48+
Mockito.when(secondEvent.isGeneratedByHistoricalExtractor()).thenReturn(true);
49+
50+
queue.directOffer(firstEvent);
51+
queue.directOffer(secondEvent);
52+
53+
Assert.assertSame(firstEvent, queue.waitedPoll());
54+
Assert.assertSame(secondEvent, queue.waitedPoll());
55+
}
56+
57+
private static class CollidingFile extends File {
58+
59+
private CollidingFile(final String pathname) {
60+
super(pathname);
61+
}
62+
63+
@Override
64+
public int hashCode() {
65+
return 0;
66+
}
67+
}
68+
}

0 commit comments

Comments
 (0)