Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -281,6 +281,7 @@ private void phaseTwoLoop(RawReader reader, MessageId to, Map<String, MessageId>
promise.complete(null);
}
});
return;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

What if some of the previous addToCompactedLedger aync write fail later, but the last id's addToCompactedLedger is complete first? Do we mark the promise as successful here?

Don't we need to wait for all outstanding addToCompactedLedger to be complete before making promise complete?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You are right, but this PR is to fix phase two can't be stopped. We can open another PR to fix this issue.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@heesung-sn I opened #21067 to fix it.

}
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,37 +19,41 @@
package org.apache.pulsar.compaction;

import static org.apache.pulsar.client.impl.RawReaderTest.extractKey;

import com.google.common.collect.Lists;
import com.google.common.collect.Sets;
import com.google.common.util.concurrent.ThreadFactoryBuilder;

import io.netty.buffer.ByteBuf;

import java.util.ArrayList;
import java.util.Enumeration;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Random;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;

import lombok.Cleanup;
import org.apache.bookkeeper.client.BookKeeper;
import org.apache.bookkeeper.client.LedgerEntry;
import org.apache.bookkeeper.client.LedgerHandle;
import org.apache.bookkeeper.mledger.Entry;
import org.apache.bookkeeper.mledger.Position;
import org.apache.bookkeeper.mledger.impl.PositionImpl;
import org.apache.pulsar.broker.ServiceConfiguration;
import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest;
import org.apache.pulsar.broker.service.persistent.PersistentTopic;
import org.apache.pulsar.client.api.MessageRoutingMode;
import org.apache.pulsar.client.api.Producer;
import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.client.api.RawMessage;
import org.apache.pulsar.client.impl.PulsarClientImpl;
import org.apache.pulsar.client.impl.RawMessageImpl;
import org.apache.pulsar.common.policies.data.ClusterData;
import org.apache.pulsar.common.protocol.Commands;
import org.apache.pulsar.common.policies.data.TenantInfoImpl;
import org.apache.pulsar.common.protocol.Commands;
import org.awaitility.Awaitility;
import org.mockito.Mockito;
import org.testng.Assert;
import org.testng.annotations.AfterMethod;
Expand All @@ -65,7 +69,6 @@ public class CompactorTest extends MockedPulsarServiceBaseTest {
protected Compactor compactor;



@BeforeMethod
@Override
public void setup() throws Exception {
Expand Down Expand Up @@ -101,16 +104,17 @@ protected Compactor getCompactor() {
return compactor;
}

private List<String> compactAndVerify(String topic, Map<String, byte[]> expected, boolean checkMetrics) throws Exception {
private List<String> compactAndVerify(String topic, Map<String, byte[]> expected, boolean checkMetrics)
throws Exception {

long compactedLedgerId = compact(topic);

LedgerHandle ledger = bk.openLedger(compactedLedgerId,
Compactor.COMPACTED_TOPIC_LEDGER_DIGEST_TYPE,
Compactor.COMPACTED_TOPIC_LEDGER_PASSWORD);
Compactor.COMPACTED_TOPIC_LEDGER_DIGEST_TYPE,
Compactor.COMPACTED_TOPIC_LEDGER_PASSWORD);
Assert.assertEquals(ledger.getLastAddConfirmed() + 1, // 0..lac
expected.size(),
"Should have as many entries as there is keys");
expected.size(),
"Should have as many entries as there is keys");

List<String> keys = new ArrayList<>();
Enumeration<LedgerEntry> entries = ledger.readEntries(0, ledger.getLastAddConfirmed());
Expand All @@ -124,7 +128,7 @@ private List<String> compactAndVerify(String topic, Map<String, byte[]> expected
byte[] bytes = new byte[payload.readableBytes()];
payload.readBytes(bytes);
Assert.assertEquals(bytes, expected.remove(key),
"Compacted version should match expected version");
"Compacted version should match expected version");
m.close();
}
if (checkMetrics) {
Expand All @@ -148,17 +152,18 @@ public void testCompaction() throws Exception {
final int numMessages = 1000;
final int maxKeys = 10;

@Cleanup
Producer<byte[]> producer = pulsarClient.newProducer().topic(topic)
.enableBatching(false)
.messageRoutingMode(MessageRoutingMode.SinglePartition)
.create();
.enableBatching(false)
.messageRoutingMode(MessageRoutingMode.SinglePartition)
.create();

Map<String, byte[]> expected = new HashMap<>();
Random r = new Random(0);

for (int j = 0; j < numMessages; j++) {
int keyIndex = r.nextInt(maxKeys);
String key = "key"+keyIndex;
String key = "key" + keyIndex;
byte[] data = ("my-message-" + key + "-" + j).getBytes();
producer.newMessage()
.key(key)
Expand All @@ -173,10 +178,11 @@ public void testCompaction() throws Exception {
public void testCompactAddCompact() throws Exception {
String topic = "persistent://my-property/use/my-ns/my-topic1";

@Cleanup
Producer<byte[]> producer = pulsarClient.newProducer().topic(topic)
.enableBatching(false)
.messageRoutingMode(MessageRoutingMode.SinglePartition)
.create();
.enableBatching(false)
.messageRoutingMode(MessageRoutingMode.SinglePartition)
.create();

Map<String, byte[]> expected = new HashMap<>();

Expand Down Expand Up @@ -210,10 +216,11 @@ public void testCompactAddCompact() throws Exception {
public void testCompactedInOrder() throws Exception {
String topic = "persistent://my-property/use/my-ns/my-topic1";

@Cleanup
Producer<byte[]> producer = pulsarClient.newProducer().topic(topic)
.enableBatching(false)
.messageRoutingMode(MessageRoutingMode.SinglePartition)
.create();
.enableBatching(false)
.messageRoutingMode(MessageRoutingMode.SinglePartition)
.create();

producer.newMessage()
.key("c")
Expand Down Expand Up @@ -256,6 +263,48 @@ public void testPhaseOneLoopTimeConfiguration() {
Assert.assertEquals(compactor.getPhaseOneLoopReadTimeoutInSeconds(), 60);
}

@Test
public void testCompactedWithConcurrentSend() throws Exception {
String topic = "persistent://my-property/use/my-ns/testCompactedWithConcurrentSend";

@Cleanup
Producer<byte[]> producer = pulsarClient.newProducer().topic(topic)
.enableBatching(false)
.messageRoutingMode(MessageRoutingMode.SinglePartition)
.create();

var future = CompletableFuture.runAsync(() -> {
for (int i = 0; i < 100; i++) {
try {
producer.newMessage().key(String.valueOf(i)).value(String.valueOf(i).getBytes()).send();
} catch (PulsarClientException e) {
throw new RuntimeException(e);
}
}
});

PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topic).get();
PulsarTopicCompactionService topicCompactionService =
(PulsarTopicCompactionService) persistentTopic.getTopicCompactionService();

Awaitility.await().untilAsserted(() -> {
long compactedLedgerId = compact(topic);
Thread.sleep(300);
Optional<CompactedTopicContext> compactedTopicContext = topicCompactionService.getCompactedTopic()
.getCompactedTopicContext();
Assert.assertTrue(compactedTopicContext.isPresent());
Assert.assertEquals(compactedTopicContext.get().ledger.getId(), compactedLedgerId);
});

Position lastCompactedPosition = topicCompactionService.getLastCompactedPosition().get();
Entry lastCompactedEntry = topicCompactionService.readLastCompactedEntry().get();

Assert.assertTrue(PositionImpl.get(lastCompactedPosition.getLedgerId(), lastCompactedPosition.getEntryId())
.compareTo(lastCompactedEntry.getLedgerId(), lastCompactedEntry.getEntryId()) >= 0);

future.join();
}

public ByteBuf extractPayload(RawMessage m) throws Exception {
ByteBuf payloadAndMetadata = m.getHeadersAndPayload();
Commands.skipChecksumIfPresent(payloadAndMetadata);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,36 +21,64 @@
import static org.apache.pulsar.compaction.Compactor.COMPACTED_TOPIC_LEDGER_PROPERTY;
import static org.apache.pulsar.compaction.Compactor.COMPACTION_SUBSCRIPTION;
import static org.testng.Assert.assertEquals;
import com.google.common.util.concurrent.ThreadFactoryBuilder;
import java.io.IOException;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import lombok.Cleanup;
import org.apache.bookkeeper.client.BookKeeper;
import org.apache.bookkeeper.mledger.Entry;
import org.apache.bookkeeper.mledger.Position;
import org.apache.bookkeeper.mledger.impl.PositionImpl;
import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest;
import org.apache.pulsar.client.admin.PulsarAdminException;
import org.apache.pulsar.client.api.MessageRoutingMode;
import org.apache.pulsar.client.api.Producer;
import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.client.impl.MessageImpl;
import org.apache.pulsar.common.policies.data.ClusterData;
import org.apache.pulsar.common.policies.data.TenantInfoImpl;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeMethod;
import org.testng.annotations.Test;

public class TopicCompactionServiceTest extends CompactorTest {
public class TopicCompactionServiceTest extends MockedPulsarServiceBaseTest {

protected ScheduledExecutorService compactionScheduler;
protected BookKeeper bk;
private TwoPhaseCompactor compactor;

@BeforeMethod
@Override
public void setup() throws Exception {
super.setup();
super.internalSetup();

admin.clusters().createCluster("test", ClusterData.builder().serviceUrl(pulsar.getWebServiceAddress()).build());
TenantInfoImpl tenantInfo = new TenantInfoImpl(Set.of("role1", "role2"), Set.of("test"));
String defaultTenant = "prop-xyz";
admin.tenants().createTenant(defaultTenant, tenantInfo);
String defaultNamespace = defaultTenant + "/ns1";
admin.namespaces().createNamespace(defaultNamespace, Set.of("test"));

compactionScheduler = Executors.newSingleThreadScheduledExecutor(
new ThreadFactoryBuilder().setNameFormat("compactor").setDaemon(true).build());
bk = pulsar.getBookKeeperClientFactory().create(
this.conf, null, null, Optional.empty(), null);
compactor = new TwoPhaseCompactor(conf, pulsarClient, bk, compactionScheduler);
}

@AfterMethod(alwaysRun = true)
@Override
public void cleanup() throws Exception {
super.internalCleanup();
bk.close();
if (compactionScheduler != null) {
compactionScheduler.shutdownNow();
}
}

@Test
Expand Down