diff --git a/pom.xml b/pom.xml index d7b5b7253d31b..209dcc14de1ee 100644 --- a/pom.xml +++ b/pom.xml @@ -235,7 +235,7 @@ flexible messaging model and an intuitive client API. 2.5.1 9+181-r4173-1 0.1.4 - 0.2 + 0.4 rename-netty-native-libs.sh diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 856f826da6eb0..358345c4cb824 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -687,12 +687,6 @@ public class ServiceConfiguration implements PulsarConfiguration { ) private String resourceUsageTransportClassName = ""; - @FieldContext( - category = CATEGORY_POLICIES, - doc = "Topic to publish usage reports to if resourceUsagePublishToTopic is enabled." - ) - private String resourceUsageTransportPublishTopicName = "non-persistent://pulsar/system/resource-usage"; - @FieldContext( dynamic = true, category = CATEGORY_POLICIES, diff --git a/pulsar-broker/pom.xml b/pulsar-broker/pom.xml index 49609c219f08e..cb130482e08c9 100644 --- a/pulsar-broker/pom.xml +++ b/pulsar-broker/pom.xml @@ -491,6 +491,8 @@ ${lightproto-maven-plugin.version} ${project.basedir}/src/main/proto/ResourceUsage.proto + generated-sources/lightproto/java + generated-sources/lightproto/java diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index 7d7787ce27e4a..ab459f1b8fa58 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -329,6 +329,11 @@ public CompletableFuture closeAsync() { } // close the service in reverse order v.s. in which they are started + if (this.resourceUsageTransportManager != null) { + this.resourceUsageTransportManager.close(); + this.resourceUsageTransportManager = null; + } + if (this.webService != null) { try { this.webService.close(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/resourcegroup/ResourceUsageTransportManager.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/resourcegroup/ResourceUsageTransportManager.java index f7064e4993c96..227bce0658ee9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/resourcegroup/ResourceUsageTransportManager.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/resourcegroup/ResourceUsageTransportManager.java @@ -66,7 +66,7 @@ private Producer createProducer() throws PulsarClientException { final int sendTimeoutSecs = 10; return pulsarClient.newProducer() - .topic(pulsarService.getConfig().getResourceUsageTransportPublishTopicName()) + .topic(RESOURCE_USAGE_TOPIC_NAME) .batchingMaxPublishDelay(publishDelayMilliSecs, TimeUnit.MILLISECONDS) .sendTimeout(sendTimeoutSecs, TimeUnit.SECONDS) .blockIfQueueFull(false) @@ -113,11 +113,12 @@ public void close() throws Exception { private class ResourceUsageReader implements ReaderListener, AutoCloseable { private final ResourceUsageInfo recdUsageInfo = new ResourceUsageInfo(); + private final Reader consumer; public ResourceUsageReader() throws PulsarClientException { consumer = pulsarClient.newReader() - .topic(pulsarService.getConfig().getResourceUsageTransportPublishTopicName()) + .topic(RESOURCE_USAGE_TOPIC_NAME) .startMessageId(MessageId.latest) .readerListener(this) .create(); @@ -130,9 +131,19 @@ public void close() throws Exception { @Override public void received(Reader reader, Message msg) { - try { - recdUsageInfo.parseFrom(Unpooled.wrappedBuffer(msg.getData()), msg.getData().length); + long publishTime = msg.getPublishTime(); + long currentTime = System.currentTimeMillis(); + long timeDelta = currentTime - publishTime; + recdUsageInfo.parseFrom(Unpooled.wrappedBuffer(msg.getData()), msg.getData().length); + if (timeDelta > TimeUnit.SECONDS.toMillis( + 2 * pulsarService.getConfig().getResourceUsageTransportPublishIntervalInSecs())) { + LOG.error("Stale resource usage msg from broker {} publish time {} current time{}", + recdUsageInfo.getBroker(), publishTime, currentTime); + staleMessageCount++; + return; + } + try { recdUsageInfo.getUsageMapsList().forEach(ru -> { ResourceUsageConsumer owner = consumerMap.get(ru.getOwner()); if (owner != null) { @@ -150,6 +161,7 @@ public void received(Reader reader, Message msg) { } private static final Logger LOG = LoggerFactory.getLogger(ResourceUsageTransportManager.class); + public static final String RESOURCE_USAGE_TOPIC_NAME = "non-persistent://pulsar/system/resource-usage"; private final PulsarService pulsarService; private final PulsarClient pulsarClient; private final ResourceUsageWriterTask pTask; @@ -159,9 +171,11 @@ public void received(Reader reader, Message msg) { private final Map consumerMap = new ConcurrentHashMap(); + private long staleMessageCount = 0; + private void createTenantAndNamespace() throws PulsarServerException, PulsarAdminException { // Create a public tenant and default namespace - TopicName topicName = TopicName.get(pulsarService.getConfig().getResourceUsageTransportPublishTopicName()); + TopicName topicName = TopicName.get(RESOURCE_USAGE_TOPIC_NAME); PulsarAdmin admin = pulsarService.getAdminClient(); ServiceConfiguration config = pulsarService.getConfig(); @@ -172,12 +186,26 @@ private void createTenantAndNamespace() throws PulsarServerException, PulsarAdmi List tenantList = admin.tenants().getTenants(); if (!tenantList.contains(tenant)) { - admin.tenants().createTenant(tenant, - new TenantInfo(Sets.newHashSet(config.getSuperUserRoles()), Sets.newHashSet(cluster))); + try { + admin.tenants().createTenant(tenant, + new TenantInfo(Sets.newHashSet(config.getSuperUserRoles()), Sets.newHashSet(cluster))); + } catch (PulsarAdminException ex1) { + if (!(ex1 instanceof PulsarAdminException.ConflictException)) { + LOG.error("Unexpected exception {} when creating tenant {}", ex1, tenant); + throw ex1; + } + } } List nsList = admin.namespaces().getNamespaces(tenant); if (!nsList.contains(namespace)) { - admin.namespaces().createNamespace(namespace); + try { + admin.namespaces().createNamespace(namespace); + } catch (PulsarAdminException ex1) { + if (!(ex1 instanceof PulsarAdminException.ConflictException)) { + LOG.error("Unexpected exception {} when creating namespace {}", ex1, namespace); + throw ex1; + } + } } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/resourcegroup/ResourceUsageTransportManagerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/resourcegroup/ResourceUsageTransportManagerTest.java index 4c1f1862a79cc..89246d8d609de 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/resourcegroup/ResourceUsageTransportManagerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/resourcegroup/ResourceUsageTransportManagerTest.java @@ -34,7 +34,6 @@ public class ResourceUsageTransportManagerTest extends MockedPulsarServiceBaseTest { - private static final String INTERNAL_TOPIC = "non-persistent://pulsar-test/test/resource-usage"; private static final int PUBLISH_INTERVAL_SECS = 1; @BeforeClass @@ -53,11 +52,10 @@ protected void cleanup() throws Exception { @Test public void testNamespaceCreation() throws Exception { ResourceUsageTransportManager tManager = new ResourceUsageTransportManager(pulsar); - TopicName topicName = TopicName.get(INTERNAL_TOPIC); + TopicName topicName = TopicName.get(ResourceUsageTransportManager.RESOURCE_USAGE_TOPIC_NAME); assertTrue(admin.tenants().getTenants().contains(topicName.getTenant())); assertTrue(admin.namespaces().getNamespaces(topicName.getTenant()).contains(topicName.getNamespace())); - } @Test @@ -116,7 +114,6 @@ public void acceptResourceUsage(String broker, ResourceUsage resourceUsage) { private void prepareData() throws PulsarAdminException { this.conf.setResourceUsageTransportClassName("org.apache.pulsar.broker.resourcegroup.ResourceUsageTransportManager"); - this.conf.setResourceUsageTransportPublishTopicName(INTERNAL_TOPIC); this.conf.setResourceUsageTransportPublishIntervalInSecs(PUBLISH_INTERVAL_SECS); admin.clusters().createCluster("test", new ClusterData(pulsar.getBrokerServiceUrl())); }