diff --git a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java index a57c7fcd1d3a..d9e7943ad806 100644 --- a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java +++ b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java @@ -135,7 +135,8 @@ abstract class AbstractSubscriptionConsumer implements AutoCloseable { private final String fileSaveDir; private final boolean fileSaveFsync; - private final Set inFlightFilesCommitContextSet = new HashSet<>(); + private final Set inFlightFilesCommitContextSet = + ConcurrentHashMap.newKeySet(); private final int thriftMaxFrameSize; private final int connectionTimeoutInMs; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/TsFileDeduplicationBlockingPendingQueue.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/TsFileDeduplicationBlockingPendingQueue.java index 5f614013e87f..14b9cb1faddc 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/TsFileDeduplicationBlockingPendingQueue.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/TsFileDeduplicationBlockingPendingQueue.java @@ -40,13 +40,13 @@ public class TsFileDeduplicationBlockingPendingQueue extends SubscriptionBlockin private static final Logger LOGGER = LoggerFactory.getLogger(TsFileDeduplicationBlockingPendingQueue.class); - private final Cache hashCodeToIsGeneratedByHistoricalExtractor; + private final Cache tsFilePathToIsGeneratedByHistoricalExtractor; public TsFileDeduplicationBlockingPendingQueue( final UnboundedBlockingPendingQueue inputPendingQueue) { super(inputPendingQueue); - this.hashCodeToIsGeneratedByHistoricalExtractor = + this.tsFilePathToIsGeneratedByHistoricalExtractor = Caffeine.newBuilder() .expireAfterAccess( SubscriptionConfig.getInstance().getSubscriptionTsFileDeduplicationWindowSeconds(), @@ -101,12 +101,13 @@ && isDuplicated((PipeTsFileInsertionEvent) sourceEvent)) { } private boolean isDuplicated(final PipeTsFileInsertionEvent event) { - final int hashCode = event.getTsFile().hashCode(); + final String tsFilePath = event.getTsFile().toPath().toAbsolutePath().normalize().toString(); final boolean isGeneratedByHistoricalExtractor = event.isGeneratedByHistoricalExtractor(); final Boolean existedIsGeneratedByHistoricalExtractor = - hashCodeToIsGeneratedByHistoricalExtractor.getIfPresent(hashCode); + tsFilePathToIsGeneratedByHistoricalExtractor.getIfPresent(tsFilePath); if (Objects.isNull(existedIsGeneratedByHistoricalExtractor)) { - hashCodeToIsGeneratedByHistoricalExtractor.put(hashCode, isGeneratedByHistoricalExtractor); + tsFilePathToIsGeneratedByHistoricalExtractor.put( + tsFilePath, isGeneratedByHistoricalExtractor); return false; } // Multiple PipeRawTabletInsertionEvents parsed from the same PipeTsFileInsertionEvent (i.e., diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/TsFileDeduplicationBlockingPendingQueueTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/TsFileDeduplicationBlockingPendingQueueTest.java new file mode 100644 index 000000000000..def98c52b359 --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/TsFileDeduplicationBlockingPendingQueueTest.java @@ -0,0 +1,68 @@ +/* + * 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.iotdb.db.subscription.broker; + +import org.apache.iotdb.commons.pipe.agent.task.connection.UnboundedBlockingPendingQueue; +import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent; +import org.apache.iotdb.db.pipe.metric.source.PipeDataRegionEventCounter; +import org.apache.iotdb.pipe.api.event.Event; + +import org.junit.Assert; +import org.junit.Test; +import org.mockito.Mockito; + +import java.io.File; + +public class TsFileDeduplicationBlockingPendingQueueTest { + + @Test + public void testDifferentPathsWithSameHashCodeAreNotDeduplicated() { + final UnboundedBlockingPendingQueue inputPendingQueue = + new UnboundedBlockingPendingQueue<>(new PipeDataRegionEventCounter()); + final TsFileDeduplicationBlockingPendingQueue queue = + new TsFileDeduplicationBlockingPendingQueue(inputPendingQueue); + + final PipeTsFileInsertionEvent firstEvent = Mockito.mock(PipeTsFileInsertionEvent.class); + Mockito.when(firstEvent.getTsFile()).thenReturn(new CollidingFile("first.tsfile")); + Mockito.when(firstEvent.isGeneratedByHistoricalExtractor()).thenReturn(false); + + final PipeTsFileInsertionEvent secondEvent = Mockito.mock(PipeTsFileInsertionEvent.class); + Mockito.when(secondEvent.getTsFile()).thenReturn(new CollidingFile("second.tsfile")); + Mockito.when(secondEvent.isGeneratedByHistoricalExtractor()).thenReturn(true); + + queue.directOffer(firstEvent); + queue.directOffer(secondEvent); + + Assert.assertSame(firstEvent, queue.waitedPoll()); + Assert.assertSame(secondEvent, queue.waitedPoll()); + } + + private static class CollidingFile extends File { + + private CollidingFile(final String pathname) { + super(pathname); + } + + @Override + public int hashCode() { + return 0; + } + } +}