From 4cce4e6cfb02dbff925d19a88480eb4f03961705 Mon Sep 17 00:00:00 2001 From: congbo Date: Fri, 23 Apr 2021 12:08:39 +0800 Subject: [PATCH 1/2] [Transaction] Fix MLTransactionLog open manageLedger name problem. --- .../service/persistent/PersistentTopic.java | 4 +++- .../stats/ManagedLedgerMetricsTest.java | 22 +++++++++++++++++++ .../impl/MLTransactionLogImpl.java | 11 ++++++---- 3 files changed, 32 insertions(+), 5 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index e1a4018580e61..755cfd3bfaa47 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -141,6 +141,7 @@ import org.apache.pulsar.compaction.CompactedTopicImpl; import org.apache.pulsar.compaction.Compactor; import org.apache.pulsar.policies.data.loadbalancer.NamespaceBundleStats; +import org.apache.pulsar.transaction.coordinator.impl.MLTransactionLogImpl; import org.apache.pulsar.utils.StatsOutputStream; import org.apache.zookeeper.KeeperException; import org.slf4j.Logger; @@ -302,7 +303,8 @@ public PersistentTopic(String topic, ManagedLedger ledger, BrokerService brokerS if (brokerService.getPulsar().getConfiguration().isTransactionCoordinatorEnabled() && !checkTopicIsEventsNames(topic) - && !topic.contains(TopicName.TRANSACTION_COORDINATOR_ASSIGN.getLocalName())) { + && !topic.contains(TopicName.TRANSACTION_COORDINATOR_ASSIGN.getLocalName()) + && !topic.contains(MLTransactionLogImpl.TRANSACTION_LOG_PREFIX)) { this.transactionBuffer = brokerService.getPulsar() .getTransactionBufferProvider().newTransactionBuffer(this, transactionCompletableFuture); } else { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ManagedLedgerMetricsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ManagedLedgerMetricsTest.java index 49176c387c3f5..9bd8561cbb25a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ManagedLedgerMetricsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ManagedLedgerMetricsTest.java @@ -22,13 +22,20 @@ import java.util.Map.Entry; import java.util.concurrent.TimeUnit; +import com.google.common.collect.Sets; +import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerMBeanImpl; import org.apache.pulsar.broker.service.BrokerTestBase; import org.apache.pulsar.broker.stats.metrics.ManagedLedgerMetrics; import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.common.naming.NamespaceName; +import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.TenantInfo; import org.apache.pulsar.common.stats.Metrics; +import org.apache.pulsar.transaction.coordinator.TransactionCoordinatorID; +import org.apache.pulsar.transaction.coordinator.impl.MLTransactionLogImpl; import org.testng.Assert; import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; @@ -87,4 +94,19 @@ public void testManagedLedgerMetrics() throws Exception { } + @Test + public void testTransactionTopic() throws Exception { + admin.tenants().createTenant(NamespaceName.SYSTEM_NAMESPACE.getTenant(), + new TenantInfo(Sets.newHashSet("appid1"), Sets.newHashSet("test"))); + admin.namespaces().createNamespace(NamespaceName.SYSTEM_NAMESPACE.toString()); + admin.topics().createPartitionedTopic(TopicName.TRANSACTION_COORDINATOR_ASSIGN.toString(), 1); + ManagedLedgerConfig managedLedgerConfig = new ManagedLedgerConfig(); + managedLedgerConfig.setMaxEntriesPerLedger(2); + new MLTransactionLogImpl(TransactionCoordinatorID.get(0), + pulsar.getManagedLedgerFactory(), managedLedgerConfig); + Thread.sleep(2000); + ManagedLedgerMetrics metrics = new ManagedLedgerMetrics(pulsar); + metrics.generate(); + } + } diff --git a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionLogImpl.java b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionLogImpl.java index 9df58dc886c47..e1a71d8be9a1c 100644 --- a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionLogImpl.java +++ b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionLogImpl.java @@ -34,6 +34,8 @@ import org.apache.pulsar.common.allocator.PulsarByteBufAllocator; import org.apache.pulsar.common.api.proto.CommandSubscribe; import org.apache.pulsar.common.naming.NamespaceName; +import org.apache.pulsar.common.naming.TopicDomain; +import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.transaction.coordinator.TransactionCoordinatorID; import org.apache.pulsar.transaction.coordinator.TransactionLog; import org.apache.pulsar.transaction.coordinator.TransactionLogReplayCallback; @@ -49,7 +51,7 @@ public class MLTransactionLogImpl implements TransactionLog { private final ManagedLedger managedLedger; - public final static String TRANSACTION_LOG_PREFIX = NamespaceName.SYSTEM_NAMESPACE + "/transaction-log-"; + public final static String TRANSACTION_LOG_PREFIX = "__transaction_log_"; private final ManagedCursor cursor; @@ -59,14 +61,15 @@ public class MLTransactionLogImpl implements TransactionLog { private final long tcId; - private final String topicName; + private final TopicName topicName; public MLTransactionLogImpl(TransactionCoordinatorID tcID, ManagedLedgerFactory managedLedgerFactory, ManagedLedgerConfig managedLedgerConfig) throws Exception { - this.topicName = TRANSACTION_LOG_PREFIX + tcID; + this.topicName = TopicName.get(TopicDomain.persistent.value(), + NamespaceName.SYSTEM_NAMESPACE, TRANSACTION_LOG_PREFIX + tcID.getId()); this.tcId = tcID.getId(); - this.managedLedger = managedLedgerFactory.open(topicName, managedLedgerConfig); + this.managedLedger = managedLedgerFactory.open(topicName.getPersistenceNamingEncoding(), managedLedgerConfig); this.cursor = managedLedger.openCursor(TRANSACTION_SUBSCRIPTION_NAME, CommandSubscribe.InitialPosition.Earliest); this.entryQueue = new SpscArrayQueue<>(2000); From 7ad85e4e8729a06aff6a28da353e969cb6b91414 Mon Sep 17 00:00:00 2001 From: congbo Date: Fri, 23 Apr 2021 13:58:07 +0800 Subject: [PATCH 2/2] Fix some comment --- .../pulsar/broker/service/persistent/PersistentTopic.java | 6 +++--- .../pulsar/broker/stats/ManagedLedgerMetricsTest.java | 1 - 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 755cfd3bfaa47..f31a55a46ed9a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -300,11 +300,11 @@ public PersistentTopic(String topic, ManagedLedger ledger, BrokerService brokerS } checkReplicatedSubscriptionControllerState(); - + TopicName topicName = TopicName.get(topic); if (brokerService.getPulsar().getConfiguration().isTransactionCoordinatorEnabled() && !checkTopicIsEventsNames(topic) - && !topic.contains(TopicName.TRANSACTION_COORDINATOR_ASSIGN.getLocalName()) - && !topic.contains(MLTransactionLogImpl.TRANSACTION_LOG_PREFIX)) { + && !topicName.getEncodedLocalName().startsWith(TopicName.TRANSACTION_COORDINATOR_ASSIGN.getLocalName()) + && !topicName.getEncodedLocalName().startsWith(MLTransactionLogImpl.TRANSACTION_LOG_PREFIX)) { this.transactionBuffer = brokerService.getPulsar() .getTransactionBufferProvider().newTransactionBuffer(this, transactionCompletableFuture); } else { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ManagedLedgerMetricsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ManagedLedgerMetricsTest.java index 9bd8561cbb25a..41ee681f8f912 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ManagedLedgerMetricsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ManagedLedgerMetricsTest.java @@ -104,7 +104,6 @@ public void testTransactionTopic() throws Exception { managedLedgerConfig.setMaxEntriesPerLedger(2); new MLTransactionLogImpl(TransactionCoordinatorID.get(0), pulsar.getManagedLedgerFactory(), managedLedgerConfig); - Thread.sleep(2000); ManagedLedgerMetrics metrics = new ManagedLedgerMetrics(pulsar); metrics.generate(); }