From e301f882c5e54c66e94803ccdcc252cdc07f69b5 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Tue, 23 Nov 2021 16:34:55 +0800 Subject: [PATCH 01/19] support entry filter --- conf/broker.conf | 3 + .../pulsar/broker/ServiceConfiguration.java | 7 ++ .../broker/service/plugin/EntryFilter.java | 62 ++++++++++++++++++ .../service/plugin/EntryFilterProvider.java | 64 +++++++++++++++++++ 4 files changed, 136 insertions(+) create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java diff --git a/conf/broker.conf b/conf/broker.conf index 700f9a5adc9b1..e79c3abb0794b 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -411,6 +411,9 @@ dispatcherReadFailureBackoffMandatoryStopTimeInMs=0 # Precise dispathcer flow control according to history message number of each entry preciseDispatcherFlowControl=false +# Class name of Pluggable entry filter that can decide whether the entry needs to be filtered +entryFilterClassName= + # Max number of concurrent lookup request broker allows to throttle heavy incoming lookup traffic maxConcurrentLookupRequest=50000 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 4407a9fbaef5c..9d93588eed426 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 @@ -829,6 +829,13 @@ public class ServiceConfiguration implements PulsarConfiguration { ) private boolean preciseDispatcherFlowControl = false; + @FieldContext( + dynamic = true, + category = CATEGORY_SERVER, + doc = " Class name of Pluggable entry filter that can decide whether the entry needs to be filtered" + ) + private String entryFilterClassName = ""; + @FieldContext( category = CATEGORY_SERVER, doc = "Whether to use streaming read dispatcher. Currently is in preview and can be changed " + diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java new file mode 100644 index 0000000000000..08d85828472ee --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java @@ -0,0 +1,62 @@ +/** + * 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.pulsar.broker.service.plugin; + +import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.ManagedCursor; +import org.apache.pulsar.broker.service.EntryBatchIndexesAcks; +import org.apache.pulsar.broker.service.EntryBatchSizes; +import org.apache.pulsar.broker.service.SendMessageInfo; +import org.apache.pulsar.broker.service.persistent.PersistentSubscription; +import org.apache.pulsar.common.api.proto.MessageMetadata; + +public interface EntryFilter { + /** + * Broker will determine whether to filter out this Entry based on the return value of this method. + * Please do not deserialize the entire Entry in this method, + * which will have a great impact on Broker's memory and CPU. + * @param entry + * @param context + * @return + */ + FilterResult filterEntry(Entry entry, FilterContext context); + + class FilterContext { + + EntryBatchSizes batchSizes; + SendMessageInfo sendMessageInfo; + EntryBatchIndexesAcks indexesAcks; + ManagedCursor cursor; + boolean isReplayRead; + PersistentSubscription subscription; + //SubscriptionOption subscriptionOption; + MessageMetadata msgMetadata; + } + + enum FilterResult { + /** + * deliver to the Consumer + */ + ACCEPT, + /** + * skip the message + */ + REJECT, + } +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java new file mode 100644 index 0000000000000..f9d07d822153e --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java @@ -0,0 +1,64 @@ +/** + * 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.pulsar.broker.service.plugin; + + +import static com.google.common.base.Preconditions.checkArgument; +import java.io.IOException; +import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.service.Subscription; +import org.apache.pulsar.transaction.coordinator.TransactionMetadataStoreProvider; + +public interface EntryFilterProvider { + /** + * Use `EntryFilterProvider` to create `EntryFilter` through reflection + * @param subscription + * @return + */ + EntryFilter createEntriesFilter(Subscription subscription) throws IOException; + + /** + * The default implementation class of EntriesFilterProvider + */ + class DefaultEntryFilterProviderImpl implements EntryFilterProvider { + private final ServiceConfiguration serviceConfiguration; + private final Subscription subscription; + + public DefaultEntryFilterProviderImpl(ServiceConfiguration serviceConfiguration, Subscription subscription) { + this.serviceConfiguration = serviceConfiguration; + this.subscription = subscription; + } + + @Override + public EntryFilter createEntriesFilter(Subscription subscription) throws IOException { + Class providerClass; + try { + providerClass = Class.forName(serviceConfiguration.getEntryFilterClassName()); + Object obj = providerClass.getDeclaredConstructor().newInstance(); + checkArgument(obj instanceof EntryFilter, + "The instance is not an instance of " + + serviceConfiguration.getEntryFilterClassName()); + return (EntryFilter) obj; + } catch (Exception e) { + throw new IOException(e); + } + } + } + +} From 3bac1bed0b6980d8081ad4a9adf0dee2548d08c7 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Wed, 24 Nov 2021 19:33:59 +0800 Subject: [PATCH 02/19] Support pluggable entry filter in Dispatcher --- conf/broker.conf | 1 + .../pulsar/broker/ServiceConfiguration.java | 3 +- .../service/AbstractBaseDispatcher.java | 35 +++++ .../persistent/PersistentSubscription.java | 3 +- .../broker/service/plugin/EntryFilter.java | 31 +++- .../service/plugin/EntryFilterProvider.java | 47 ++---- .../service/plugin/EntryFilterForTest.java | 56 +++++++ .../service/plugin/FilterEntryTest.java | 138 ++++++++++++++++++ 8 files changed, 268 insertions(+), 46 deletions(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java diff --git a/conf/broker.conf b/conf/broker.conf index e79c3abb0794b..34a647561d1b4 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -412,6 +412,7 @@ dispatcherReadFailureBackoffMandatoryStopTimeInMs=0 preciseDispatcherFlowControl=false # Class name of Pluggable entry filter that can decide whether the entry needs to be filtered +# You can use this class to decide which entries can be sent to consumers. entryFilterClassName= # Max number of concurrent lookup request broker allows to throttle heavy incoming lookup traffic 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 9d93588eed426..94d4ffb3ba142 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 @@ -832,7 +832,8 @@ public class ServiceConfiguration implements PulsarConfiguration { @FieldContext( dynamic = true, category = CATEGORY_SERVER, - doc = " Class name of Pluggable entry filter that can decide whether the entry needs to be filtered" + doc = " Class name of Pluggable entry filter that can decide whether the entry needs to be filtered." + + "You can use this class to decide which entries can be sent to consumers." ) private String entryFilterClassName = ""; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index b53de2f791ba4..bb85e081d9289 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -27,10 +27,13 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.intercept.BrokerInterceptor; import org.apache.pulsar.broker.service.persistent.PersistentTopic; +import org.apache.pulsar.broker.service.plugin.EntryFilter; +import org.apache.pulsar.broker.service.plugin.EntryFilterProvider; import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.common.api.proto.CommandAck.AckType; import org.apache.pulsar.common.api.proto.MessageMetadata; @@ -48,11 +51,14 @@ public abstract class AbstractBaseDispatcher implements Dispatcher { protected final ServiceConfiguration serviceConfig; protected final boolean dispatchThrottlingOnBatchMessageEnabled; + protected final EntryFilter entryFilter; protected AbstractBaseDispatcher(Subscription subscription, ServiceConfiguration serviceConfig) { this.subscription = subscription; this.serviceConfig = serviceConfig; this.dispatchThrottlingOnBatchMessageEnabled = serviceConfig.isDispatchThrottlingOnBatchMessageEnabled(); + this.entryFilter = StringUtils.isNotBlank(serviceConfig.getEntryFilterClassName()) + ? EntryFilterProvider.createEntryFilter(serviceConfig.getEntryFilterClassName()) : null; } /** @@ -113,6 +119,7 @@ public int filterEntriesForConsumer(Optional entryWrapper, int e long totalBytes = 0; int totalChunkedMessages = 0; int totalEntries = 0; + EntryFilter.FilterContext filterContext = new EntryFilter.FilterContext(); for (int i = 0, entriesSize = entries.size(); i < entriesSize; i++) { Entry entry = entries.get(i); if (entry == null) { @@ -127,6 +134,19 @@ public int filterEntriesForConsumer(Optional entryWrapper, int e msgMetadata = msgMetadata == null ? Commands.peekMessageMetadata(metadataAndPayload, subscription.toString(), -1) : msgMetadata; + if (entryFilter != null) { + fillContext(filterContext, batchSizes, sendMessageInfo, indexesAcks, cursor, isReplayRead, + msgMetadata, subscription); + EntryFilter.FilterResult result = entryFilter.filterEntry(entry, filterContext); + if (EntryFilter.FilterResult.REJECT == result) { + PositionImpl pos = (PositionImpl) entry.getPosition(); + entries.set(i, null); + entry.release(); + subscription.acknowledgeMessage(Collections.singletonList(pos), AckType.Individual, + Collections.emptyMap()); + continue; + } + } if (!isReplayRead && msgMetadata != null && msgMetadata.hasTxnidMostBits() && msgMetadata.hasTxnidLeastBits()) { if (Markers.isTxnMarker(msgMetadata)) { @@ -189,6 +209,21 @@ && trackDelayedDelivery(entry.getLedgerId(), entry.getEntryId(), msgMetadata)) { return totalEntries; } + private void fillContext(EntryFilter.FilterContext context, + EntryBatchSizes batchSizes, SendMessageInfo sendMessageInfo, + EntryBatchIndexesAcks indexesAcks, ManagedCursor cursor, + boolean isReplayRead, MessageMetadata msgMetadata, + Subscription subscription) { + context.reset(); + context.setBatchSizes(batchSizes); + context.setSendMessageInfo(sendMessageInfo); + context.setIndexesAcks(indexesAcks); + context.setCursor(cursor); + context.setReplayRead(isReplayRead); + context.setMsgMetadata(msgMetadata); + context.setSubscription(subscription); + } + /** * Determine whether the number of consumers on the subscription reaches the threshold. * @return diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java index fbe11fbdf0dc1..8c8dbdb2e75ef 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java @@ -48,6 +48,7 @@ import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.commons.collections4.MapUtils; import org.apache.commons.lang3.tuple.MutablePair; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.intercept.BrokerInterceptor; @@ -158,7 +159,7 @@ public PersistentSubscription(PersistentTopic topic, String subscriptionName, Ma this.fullName = MoreObjects.toStringHelper(this).add("topic", topicName).add("name", subName).toString(); this.expiryMonitor = new PersistentMessageExpiryMonitor(topicName, subscriptionName, cursor, this); this.setReplicated(replicated); - this.subscriptionProperties = subscriptionProperties == null + this.subscriptionProperties = MapUtils.isEmpty(subscriptionProperties) ? new HashMap<>() : Collections.unmodifiableMap(subscriptionProperties); if (topic.getBrokerService().getPulsar().getConfig().isTransactionCoordinatorEnabled() && !checkTopicIsEventsNames(TopicName.get(topicName))) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java index 08d85828472ee..893ac0c4f6489 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java @@ -18,11 +18,14 @@ */ package org.apache.pulsar.broker.service.plugin; +import lombok.Data; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.pulsar.broker.service.EntryBatchIndexesAcks; import org.apache.pulsar.broker.service.EntryBatchSizes; import org.apache.pulsar.broker.service.SendMessageInfo; +import org.apache.pulsar.broker.service.Subscription; +import org.apache.pulsar.broker.service.SubscriptionOption; import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.common.api.proto.MessageMetadata; @@ -37,16 +40,28 @@ public interface EntryFilter { */ FilterResult filterEntry(Entry entry, FilterContext context); + @Data class FilterContext { - EntryBatchSizes batchSizes; - SendMessageInfo sendMessageInfo; - EntryBatchIndexesAcks indexesAcks; - ManagedCursor cursor; - boolean isReplayRead; - PersistentSubscription subscription; - //SubscriptionOption subscriptionOption; - MessageMetadata msgMetadata; + private EntryBatchSizes batchSizes; + private SendMessageInfo sendMessageInfo; + private EntryBatchIndexesAcks indexesAcks; + private ManagedCursor cursor; + private boolean isReplayRead; + private Subscription subscription; + private SubscriptionOption subscriptionOption; + private MessageMetadata msgMetadata; + + public void reset() { + batchSizes = null; + sendMessageInfo = null; + indexesAcks = null; + cursor = null; + isReplayRead = false; + subscription = null; + subscriptionOption = null; + msgMetadata = null; + } } enum FilterResult { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java index f9d07d822153e..2c7b5f9439fd3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java @@ -18,47 +18,22 @@ */ package org.apache.pulsar.broker.service.plugin; - import static com.google.common.base.Preconditions.checkArgument; -import java.io.IOException; -import org.apache.pulsar.broker.ServiceConfiguration; -import org.apache.pulsar.broker.service.Subscription; -import org.apache.pulsar.transaction.coordinator.TransactionMetadataStoreProvider; -public interface EntryFilterProvider { - /** - * Use `EntryFilterProvider` to create `EntryFilter` through reflection - * @param subscription - * @return - */ - EntryFilter createEntriesFilter(Subscription subscription) throws IOException; +public class EntryFilterProvider { /** - * The default implementation class of EntriesFilterProvider + * create entry filter instance */ - class DefaultEntryFilterProviderImpl implements EntryFilterProvider { - private final ServiceConfiguration serviceConfiguration; - private final Subscription subscription; - - public DefaultEntryFilterProviderImpl(ServiceConfiguration serviceConfiguration, Subscription subscription) { - this.serviceConfiguration = serviceConfiguration; - this.subscription = subscription; - } - - @Override - public EntryFilter createEntriesFilter(Subscription subscription) throws IOException { - Class providerClass; - try { - providerClass = Class.forName(serviceConfiguration.getEntryFilterClassName()); - Object obj = providerClass.getDeclaredConstructor().newInstance(); - checkArgument(obj instanceof EntryFilter, - "The instance is not an instance of " - + serviceConfiguration.getEntryFilterClassName()); - return (EntryFilter) obj; - } catch (Exception e) { - throw new IOException(e); - } + public static EntryFilter createEntryFilter(String className) { + Class entryFilterClass; + try { + entryFilterClass = Class.forName(className); + Object obj = entryFilterClass.getDeclaredConstructor().newInstance(); + checkArgument(obj instanceof EntryFilter, "The instance is not an instance of " + className); + return (EntryFilter) obj; + } catch (Exception e) { + throw new RuntimeException(e); } } - } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java new file mode 100644 index 0000000000000..5293849f73ad1 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java @@ -0,0 +1,56 @@ +/** + * 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.pulsar.broker.service.plugin; + + +import java.util.List; +import org.apache.bookkeeper.mledger.Entry; +import org.apache.commons.collections4.MapUtils; +import org.apache.pulsar.broker.service.persistent.PersistentSubscription; +import org.apache.pulsar.common.api.proto.KeyValue; + +public class EntryFilterForTest implements EntryFilter { + @Override + public FilterResult filterEntry(Entry entry, FilterContext context) { + if (context.getMsgMetadata() == null || context.getMsgMetadata().getPropertiesCount() <= 0) { + return FilterResult.ACCEPT; + } + List list = context.getMsgMetadata().getPropertiesList(); + // filter by subscription properties + PersistentSubscription subscription = (PersistentSubscription) context.getSubscription(); + if (!MapUtils.isEmpty(subscription.getSubscriptionProperties())) { + for (KeyValue keyValue : list) { + if(subscription.getSubscriptionProperties().containsKey(keyValue.getKey())){ + System.out.println("1111111111111111" + keyValue.getKey()); + return FilterResult.ACCEPT; + } + } + return FilterResult.REJECT; + } + // filter by string + for (KeyValue keyValue : list) { + if ("ACCEPT".equalsIgnoreCase(keyValue.getKey())) { + return FilterResult.ACCEPT; + } else if ("REJECT".equalsIgnoreCase(keyValue.getKey())){ + return FilterResult.REJECT; + } + } + return null; + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java new file mode 100644 index 0000000000000..d46aac1d722ed --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java @@ -0,0 +1,138 @@ +/** + * 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.pulsar.broker.service.plugin; + +import static org.testng.AssertJUnit.assertEquals; +import static org.testng.AssertJUnit.assertNotNull; +import java.util.HashMap; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.TimeUnit; +import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.pulsar.broker.service.BrokerTestBase; +import org.apache.pulsar.broker.service.persistent.PersistentSubscription; +import org.apache.pulsar.client.api.Consumer; +import org.apache.pulsar.client.api.Message; +import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.impl.MessageIdImpl; +import org.awaitility.Awaitility; +import org.testng.annotations.AfterMethod; +import org.testng.annotations.BeforeMethod; +import org.testng.annotations.Test; + +@Test(groups = "broker") +public class FilterEntryTest extends BrokerTestBase { + @BeforeMethod + @Override + protected void setup() throws Exception { + baseSetup(); + } + + @AfterMethod + @Override + protected void cleanup() throws Exception { + internalCleanup(); + } + + public void testFilter() throws Exception { + internalCleanup(); + conf.setEntryFilterClassName("org.apache.pulsar.broker.service.plugin.EntryFilterForTest"); + baseSetup(); + + String topic = "persistent://prop/ns-abc/topic" + UUID.randomUUID(); + String subName = "sub"; + Consumer consumer = pulsarClient.newConsumer(Schema.STRING).topic(topic) + .subscriptionName(subName).subscribe(); + Producer producer = pulsarClient.newProducer(Schema.STRING) + .enableBatching(false) + .topic(topic).create(); + for (int i = 0; i < 10; i++) { + producer.send("test"); + } + + int counter = 0; + while (true) { + Message message = consumer.receive(1, TimeUnit.SECONDS); + if (message != null) { + counter++; + consumer.acknowledge(message); + } else { + break; + } + } + // All normal messages can be received + assertEquals(10, counter); + MessageIdImpl lastMsgId = null; + for (int i = 0; i < 10; i++) { + lastMsgId = (MessageIdImpl) producer.newMessage().property("REJECT", "").value("1").send(); + } + counter = 0; + while (true) { + Message message = consumer.receive(1, TimeUnit.SECONDS); + if (message != null) { + counter++; + consumer.acknowledge(message); + } else { + break; + } + } + // REJECT messages are filtered out + assertEquals(0, counter); + + // All messages should be acked, check the MarkDeletedPosition + PersistentSubscription subscription = + (PersistentSubscription) pulsar.getBrokerService() + .getTopicReference(topic).get().getSubscription(subName); + assertNotNull(lastMsgId); + MessageIdImpl finalLastMsgId = lastMsgId; + Awaitility.await().untilAsserted(() -> { + PositionImpl position = (PositionImpl) subscription.getCursor().getMarkDeletedPosition(); + assertEquals(position.getLedgerId(), finalLastMsgId.getLedgerId()); + assertEquals(position.getEntryId(), finalLastMsgId.getEntryId()); + }); + consumer.close(); + + Map map = new HashMap<>(); + map.put("1","1"); + map.put("2","2"); + consumer = pulsarClient.newConsumer(Schema.STRING).topic(topic).subscriptionProperties(map) + .subscriptionName(subName).subscribe(); + for (int i = 0; i < 10; i++) { + producer.newMessage().property(String.valueOf(i), String.valueOf(i)).value("1").send(); + } + counter = 0; + while (true) { + Message message = consumer.receive(1, TimeUnit.SECONDS); + if (message != null) { + counter++; + consumer.acknowledge(message); + } else { + break; + } + } + assertEquals(2, counter); + + producer.close(); + consumer.close(); + + } + + +} From b1b033a14f583ab8042cb06d82dc5900af0e1cbc Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Thu, 25 Nov 2021 13:06:02 +0800 Subject: [PATCH 03/19] fix check style --- .../pulsar/broker/service/AbstractBaseDispatcher.java | 2 -- .../pulsar/broker/service/{plugin => }/EntryFilter.java | 8 +------- .../broker/service/{plugin => }/EntryFilterProvider.java | 2 +- .../pulsar/broker/service/plugin/EntryFilterForTest.java | 1 + 4 files changed, 3 insertions(+), 10 deletions(-) rename pulsar-broker/src/main/java/org/apache/pulsar/broker/service/{plugin => }/EntryFilter.java (84%) rename pulsar-broker/src/main/java/org/apache/pulsar/broker/service/{plugin => }/EntryFilterProvider.java (96%) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index bb85e081d9289..e26f92699e23d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -32,8 +32,6 @@ import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.intercept.BrokerInterceptor; import org.apache.pulsar.broker.service.persistent.PersistentTopic; -import org.apache.pulsar.broker.service.plugin.EntryFilter; -import org.apache.pulsar.broker.service.plugin.EntryFilterProvider; import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.common.api.proto.CommandAck.AckType; import org.apache.pulsar.common.api.proto.MessageMetadata; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java similarity index 84% rename from pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java rename to pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java index 893ac0c4f6489..a75240cf817aa 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java @@ -16,17 +16,11 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.pulsar.broker.service.plugin; +package org.apache.pulsar.broker.service; import lombok.Data; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; -import org.apache.pulsar.broker.service.EntryBatchIndexesAcks; -import org.apache.pulsar.broker.service.EntryBatchSizes; -import org.apache.pulsar.broker.service.SendMessageInfo; -import org.apache.pulsar.broker.service.Subscription; -import org.apache.pulsar.broker.service.SubscriptionOption; -import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.common.api.proto.MessageMetadata; public interface EntryFilter { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java similarity index 96% rename from pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java rename to pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java index 2c7b5f9439fd3..06c27fde61da4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java @@ -16,7 +16,7 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.pulsar.broker.service.plugin; +package org.apache.pulsar.broker.service; import static com.google.common.base.Preconditions.checkArgument; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java index 5293849f73ad1..6ecd1d1d6faa4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java @@ -22,6 +22,7 @@ import java.util.List; import org.apache.bookkeeper.mledger.Entry; import org.apache.commons.collections4.MapUtils; +import org.apache.pulsar.broker.service.EntryFilter; import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.common.api.proto.KeyValue; From 4dccf98799d0ee874e6bf8421127861a89754845 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Thu, 25 Nov 2021 13:50:42 +0800 Subject: [PATCH 04/19] Support multi entry filter --- conf/broker.conf | 3 +- .../pulsar/broker/ServiceConfiguration.java | 3 +- .../service/AbstractBaseDispatcher.java | 22 ++++++--- .../pulsar/broker/service/EntryFilter.java | 4 +- .../broker/service/EntryFilterProvider.java | 2 +- .../service/plugin/EntryFilterForTest.java | 13 ----- .../service/plugin/EntryFilterForTest2.java | 47 +++++++++++++++++++ .../service/plugin/FilterEntryTest.java | 4 +- 8 files changed, 73 insertions(+), 25 deletions(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest2.java diff --git a/conf/broker.conf b/conf/broker.conf index 34a647561d1b4..c7afe10dc8b8b 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -413,7 +413,8 @@ preciseDispatcherFlowControl=false # Class name of Pluggable entry filter that can decide whether the entry needs to be filtered # You can use this class to decide which entries can be sent to consumers. -entryFilterClassName= +# Multiple classes need to be separated by commas. +entryFilterClassNames= # Max number of concurrent lookup request broker allows to throttle heavy incoming lookup traffic maxConcurrentLookupRequest=50000 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 94d4ffb3ba142..5e57cbb5038a4 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 @@ -834,8 +834,9 @@ public class ServiceConfiguration implements PulsarConfiguration { category = CATEGORY_SERVER, doc = " Class name of Pluggable entry filter that can decide whether the entry needs to be filtered." + "You can use this class to decide which entries can be sent to consumers." + + "Multiple classes need to be separated by commas." ) - private String entryFilterClassName = ""; + private List entryFilterClassNames = new ArrayList<>(); @FieldContext( category = CATEGORY_SERVER, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index e26f92699e23d..1e95de6438e89 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -20,6 +20,7 @@ package org.apache.pulsar.broker.service; import io.netty.buffer.ByteBuf; +import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Optional; @@ -27,7 +28,7 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.impl.PositionImpl; -import org.apache.commons.lang3.StringUtils; +import org.apache.commons.collections4.CollectionUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.intercept.BrokerInterceptor; @@ -49,14 +50,17 @@ public abstract class AbstractBaseDispatcher implements Dispatcher { protected final ServiceConfiguration serviceConfig; protected final boolean dispatchThrottlingOnBatchMessageEnabled; - protected final EntryFilter entryFilter; + protected final List entryFilters = new ArrayList<>(); protected AbstractBaseDispatcher(Subscription subscription, ServiceConfiguration serviceConfig) { this.subscription = subscription; this.serviceConfig = serviceConfig; this.dispatchThrottlingOnBatchMessageEnabled = serviceConfig.isDispatchThrottlingOnBatchMessageEnabled(); - this.entryFilter = StringUtils.isNotBlank(serviceConfig.getEntryFilterClassName()) - ? EntryFilterProvider.createEntryFilter(serviceConfig.getEntryFilterClassName()) : null; + if (CollectionUtils.isNotEmpty(serviceConfig.getEntryFilterClassNames())) { + for (String entryFilterClassName : serviceConfig.getEntryFilterClassNames()) { + entryFilters.add(EntryFilterProvider.createEntryFilter(entryFilterClassName)); + } + } } /** @@ -132,10 +136,16 @@ public int filterEntriesForConsumer(Optional entryWrapper, int e msgMetadata = msgMetadata == null ? Commands.peekMessageMetadata(metadataAndPayload, subscription.toString(), -1) : msgMetadata; - if (entryFilter != null) { + if (CollectionUtils.isNotEmpty(entryFilters)) { fillContext(filterContext, batchSizes, sendMessageInfo, indexesAcks, cursor, isReplayRead, msgMetadata, subscription); - EntryFilter.FilterResult result = entryFilter.filterEntry(entry, filterContext); + EntryFilter.FilterResult result = EntryFilter.FilterResult.REJECT; + for (EntryFilter entryFilter : entryFilters) { + if (entryFilter.filterEntry(entry, filterContext) == EntryFilter.FilterResult.ACCEPT) { + result = EntryFilter.FilterResult.ACCEPT; + break; + } + } if (EntryFilter.FilterResult.REJECT == result) { PositionImpl pos = (PositionImpl) entry.getPosition(); entries.set(i, null); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java index a75240cf817aa..4fef6d759354e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java @@ -60,11 +60,11 @@ public void reset() { enum FilterResult { /** - * deliver to the Consumer + * deliver to the Consumer. */ ACCEPT, /** - * skip the message + * skip the message. */ REJECT, } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java index 06c27fde61da4..e258b7fbef48f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java @@ -23,7 +23,7 @@ public class EntryFilterProvider { /** - * create entry filter instance + * create entry filter instance. */ public static EntryFilter createEntryFilter(String className) { Class entryFilterClass; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java index 6ecd1d1d6faa4..e5640e5c7f3c8 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java @@ -21,9 +21,7 @@ import java.util.List; import org.apache.bookkeeper.mledger.Entry; -import org.apache.commons.collections4.MapUtils; import org.apache.pulsar.broker.service.EntryFilter; -import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.common.api.proto.KeyValue; public class EntryFilterForTest implements EntryFilter { @@ -33,17 +31,6 @@ public FilterResult filterEntry(Entry entry, FilterContext context) { return FilterResult.ACCEPT; } List list = context.getMsgMetadata().getPropertiesList(); - // filter by subscription properties - PersistentSubscription subscription = (PersistentSubscription) context.getSubscription(); - if (!MapUtils.isEmpty(subscription.getSubscriptionProperties())) { - for (KeyValue keyValue : list) { - if(subscription.getSubscriptionProperties().containsKey(keyValue.getKey())){ - System.out.println("1111111111111111" + keyValue.getKey()); - return FilterResult.ACCEPT; - } - } - return FilterResult.REJECT; - } // filter by string for (KeyValue keyValue : list) { if ("ACCEPT".equalsIgnoreCase(keyValue.getKey())) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest2.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest2.java new file mode 100644 index 0000000000000..b6ab2afb5fe5b --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest2.java @@ -0,0 +1,47 @@ +/** + * 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.pulsar.broker.service.plugin; + + +import java.util.List; +import org.apache.bookkeeper.mledger.Entry; +import org.apache.commons.collections4.MapUtils; +import org.apache.pulsar.broker.service.EntryFilter; +import org.apache.pulsar.broker.service.persistent.PersistentSubscription; +import org.apache.pulsar.common.api.proto.KeyValue; + +public class EntryFilterForTest2 implements EntryFilter { + @Override + public FilterResult filterEntry(Entry entry, FilterContext context) { + if (context.getMsgMetadata() == null || context.getMsgMetadata().getPropertiesCount() <= 0) { + return FilterResult.ACCEPT; + } + List list = context.getMsgMetadata().getPropertiesList(); + // filter by subscription properties + PersistentSubscription subscription = (PersistentSubscription) context.getSubscription(); + if (!MapUtils.isEmpty(subscription.getSubscriptionProperties())) { + for (KeyValue keyValue : list) { + if(subscription.getSubscriptionProperties().containsKey(keyValue.getKey())){ + return FilterResult.ACCEPT; + } + } + } + return FilterResult.REJECT; + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java index d46aac1d722ed..c72cabfa73221 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java @@ -20,6 +20,7 @@ import static org.testng.AssertJUnit.assertEquals; import static org.testng.AssertJUnit.assertNotNull; +import java.util.Arrays; import java.util.HashMap; import java.util.Map; import java.util.UUID; @@ -53,7 +54,8 @@ protected void cleanup() throws Exception { public void testFilter() throws Exception { internalCleanup(); - conf.setEntryFilterClassName("org.apache.pulsar.broker.service.plugin.EntryFilterForTest"); + conf.setEntryFilterClassNames(Arrays.asList("org.apache.pulsar.broker.service.plugin.EntryFilterForTest", + "org.apache.pulsar.broker.service.plugin.EntryFilterForTest2")); baseSetup(); String topic = "persistent://prop/ns-abc/topic" + UUID.randomUUID(); From a215b6bec2199724545a14c7da1adceb1f00612a Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Thu, 25 Nov 2021 17:20:33 +0800 Subject: [PATCH 05/19] Update pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java Co-authored-by: Anonymitaet <50226895+Anonymitaet@users.noreply.github.com> --- .../main/java/org/apache/pulsar/broker/service/EntryFilter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java index 4fef6d759354e..276f1190d3fef 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java @@ -25,7 +25,7 @@ public interface EntryFilter { /** - * Broker will determine whether to filter out this Entry based on the return value of this method. + * Broker determines whether to filter out this entry based on the return value of this method. * Please do not deserialize the entire Entry in this method, * which will have a great impact on Broker's memory and CPU. * @param entry From 942812f0b798286b20605f28c6f25e58ad5dd911 Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Thu, 25 Nov 2021 17:20:40 +0800 Subject: [PATCH 06/19] Update pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java Co-authored-by: Anonymitaet <50226895+Anonymitaet@users.noreply.github.com> --- .../main/java/org/apache/pulsar/broker/service/EntryFilter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java index 276f1190d3fef..4a71142e3dcae 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java @@ -27,7 +27,7 @@ public interface EntryFilter { /** * Broker determines whether to filter out this entry based on the return value of this method. * Please do not deserialize the entire Entry in this method, - * which will have a great impact on Broker's memory and CPU. + * which has a great impact on the broker's memory and CPU. * @param entry * @param context * @return From e49760f9fe80515047802c0b3e6b1ed660055519 Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Thu, 25 Nov 2021 17:20:47 +0800 Subject: [PATCH 07/19] Update pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java Co-authored-by: Anonymitaet <50226895+Anonymitaet@users.noreply.github.com> --- .../main/java/org/apache/pulsar/broker/service/EntryFilter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java index 4a71142e3dcae..572566d5ae1bb 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java @@ -60,7 +60,7 @@ public void reset() { enum FilterResult { /** - * deliver to the Consumer. + * deliver to the consumer. */ ACCEPT, /** From 1c0657032d286db969e13066235c3bf978f41ae7 Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Thu, 25 Nov 2021 17:21:01 +0800 Subject: [PATCH 08/19] Update pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java Co-authored-by: Anonymitaet <50226895+Anonymitaet@users.noreply.github.com> --- .../main/java/org/apache/pulsar/broker/service/EntryFilter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java index 572566d5ae1bb..9bb11cb596976 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java @@ -26,7 +26,7 @@ public interface EntryFilter { /** * Broker determines whether to filter out this entry based on the return value of this method. - * Please do not deserialize the entire Entry in this method, + * Do not deserialize the entire entry in this method, * which has a great impact on the broker's memory and CPU. * @param entry * @param context From 7e9f5b5106224088f672ff583c3ecad8260aa987 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Fri, 26 Nov 2021 13:05:43 +0800 Subject: [PATCH 09/19] Address comment --- conf/broker.conf | 5 +- .../pulsar/broker/ServiceConfiguration.java | 13 +- .../service/AbstractBaseDispatcher.java | 36 +++-- .../pulsar/broker/service/AbstractTopic.java | 1 + .../pulsar/broker/service/BrokerService.java | 5 + .../pulsar/broker/service/EntryFilter.java | 61 ++++++++ .../broker/service/EntryFilterProvider.java | 133 ++++++++++++++++-- .../apache/pulsar/broker/service/Topic.java | 6 + .../service/plugin/FilterEntryTest.java | 20 +-- 9 files changed, 246 insertions(+), 34 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index c7afe10dc8b8b..496af67b76faf 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -414,7 +414,10 @@ preciseDispatcherFlowControl=false # Class name of Pluggable entry filter that can decide whether the entry needs to be filtered # You can use this class to decide which entries can be sent to consumers. # Multiple classes need to be separated by commas. -entryFilterClassNames= +entryFilterNames= + +# The directory for all the entry filter implementations +entryFiltersDirectory= # Max number of concurrent lookup request broker allows to throttle heavy incoming lookup traffic maxConcurrentLookupRequest=50000 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 5e57cbb5038a4..ec9304324e7b0 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 @@ -832,11 +832,18 @@ public class ServiceConfiguration implements PulsarConfiguration { @FieldContext( dynamic = true, category = CATEGORY_SERVER, - doc = " Class name of Pluggable entry filter that can decide whether the entry needs to be filtered." + doc = " Class name of pluggable entry filter that decides whether the entry needs to be filtered." + "You can use this class to decide which entries can be sent to consumers." - + "Multiple classes need to be separated by commas." + + "Multiple names need to be separated by commas." ) - private List entryFilterClassNames = new ArrayList<>(); + private List entryFilterNames = new ArrayList<>(); + + @FieldContext( + dynamic = true, + category = CATEGORY_SERVER, + doc = " The directory for all the entry filter implementations." + ) + private String entryFiltersDirectory = ""; @FieldContext( category = CATEGORY_SERVER, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index 1e95de6438e89..519612fe33f33 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -19,8 +19,8 @@ package org.apache.pulsar.broker.service; +import com.google.common.collect.ImmutableList; import io.netty.buffer.ByteBuf; -import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Optional; @@ -29,6 +29,7 @@ import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.commons.collections4.CollectionUtils; +import org.apache.commons.collections4.MapUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.intercept.BrokerInterceptor; @@ -50,16 +51,20 @@ public abstract class AbstractBaseDispatcher implements Dispatcher { protected final ServiceConfiguration serviceConfig; protected final boolean dispatchThrottlingOnBatchMessageEnabled; - protected final List entryFilters = new ArrayList<>(); + /** + * Entry filters in Broker. + * Not set to final, for the convenience of testing mock. + */ + protected ImmutableList entryFilters; protected AbstractBaseDispatcher(Subscription subscription, ServiceConfiguration serviceConfig) { this.subscription = subscription; this.serviceConfig = serviceConfig; this.dispatchThrottlingOnBatchMessageEnabled = serviceConfig.isDispatchThrottlingOnBatchMessageEnabled(); - if (CollectionUtils.isNotEmpty(serviceConfig.getEntryFilterClassNames())) { - for (String entryFilterClassName : serviceConfig.getEntryFilterClassNames()) { - entryFilters.add(EntryFilterProvider.createEntryFilter(entryFilterClassName)); - } + if (MapUtils.isNotEmpty(subscription.getTopic().getBrokerService().getEntryFilters())) { + this.entryFilters = subscription.getTopic().getBrokerService().getEntryFilters().values().asList(); + } else { + this.entryFilters = ImmutableList.of(); } } @@ -139,13 +144,7 @@ public int filterEntriesForConsumer(Optional entryWrapper, int e if (CollectionUtils.isNotEmpty(entryFilters)) { fillContext(filterContext, batchSizes, sendMessageInfo, indexesAcks, cursor, isReplayRead, msgMetadata, subscription); - EntryFilter.FilterResult result = EntryFilter.FilterResult.REJECT; - for (EntryFilter entryFilter : entryFilters) { - if (entryFilter.filterEntry(entry, filterContext) == EntryFilter.FilterResult.ACCEPT) { - result = EntryFilter.FilterResult.ACCEPT; - break; - } - } + EntryFilter.FilterResult result = getFilterResult(filterContext, entry); if (EntryFilter.FilterResult.REJECT == result) { PositionImpl pos = (PositionImpl) entry.getPosition(); entries.set(i, null); @@ -217,6 +216,17 @@ && trackDelayedDelivery(entry.getLedgerId(), entry.getEntryId(), msgMetadata)) { return totalEntries; } + private EntryFilter.FilterResult getFilterResult(EntryFilter.FilterContext filterContext, Entry entry) { + EntryFilter.FilterResult result = EntryFilter.FilterResult.REJECT; + for (EntryFilter entryFilter : entryFilters) { + if (entryFilter.filterEntry(entry, filterContext) == EntryFilter.FilterResult.ACCEPT) { + result = EntryFilter.FilterResult.ACCEPT; + break; + } + } + return result; + } + private void fillContext(EntryFilter.FilterContext context, EntryBatchSizes batchSizes, SendMessageInfo sendMessageInfo, EntryBatchIndexesAcks indexesAcks, ManagedCursor cursor, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java index 26c591d31c282..ac83775605d51 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java @@ -318,6 +318,7 @@ public Map getProducers() { } + @Override public BrokerService getBrokerService() { return brokerService; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index e6a25742b8df6..2f8f7ab7cc676 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -26,6 +26,7 @@ import static org.apache.pulsar.broker.PulsarService.isTransactionSystemTopic; import static org.apache.pulsar.common.events.EventsTopicNames.checkTopicIsEventsNames; import com.google.common.annotations.VisibleForTesting; +import com.google.common.collect.ImmutableMap; import com.google.common.collect.Lists; import com.google.common.collect.Maps; import com.google.common.collect.Queues; @@ -267,6 +268,7 @@ public class BrokerService implements Closeable { private boolean preciseTopicPublishRateLimitingEnable; private final LongAdder pausedConnections = new LongAdder(); private BrokerInterceptor interceptor; + private ImmutableMap entryFilters; private Set brokerEntryMetadataInterceptors; @@ -299,6 +301,9 @@ public BrokerService(PulsarService pulsar, EventLoopGroup eventLoopGroup) throws .newSingleThreadScheduledExecutor(new DefaultThreadFactory("pulsar-stats-updater")); this.authorizationService = new AuthorizationService( pulsar.getConfiguration(), pulsar().getPulsarResources()); + if (!pulsar.getConfiguration().getEntryFilterNames().isEmpty()) { + this.entryFilters = EntryFilterProvider.createEntryFilters(pulsar.getConfiguration()); + } pulsar.getLocalMetadataStore().registerListener(this::handleMetadataChanges); pulsar.getConfigurationMetadataStore().registerListener(this::handleMetadataChanges); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java index 4fef6d759354e..4dfcd56773031 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java @@ -18,12 +18,18 @@ */ package org.apache.pulsar.broker.service; +import java.nio.file.Path; +import java.util.Map; +import java.util.TreeMap; import lombok.Data; +import lombok.NoArgsConstructor; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.pulsar.common.api.proto.MessageMetadata; +import org.apache.pulsar.common.nar.NarClassLoader; public interface EntryFilter { + /** * Broker will determine whether to filter out this Entry based on the return value of this method. * Please do not deserialize the entire Entry in this method, @@ -68,4 +74,59 @@ enum FilterResult { */ REJECT, } + + @Data + class EntryFilterDefinitions { + private final Map filters = new TreeMap<>(); + } + + @Data + @NoArgsConstructor + class EntryFilterMetaData { + + /** + * The definition of the broker interceptor. + */ + private EntryFilterDefinition definition; + + /** + * The path to the handler package. + */ + private Path archivePath; + } + + @Data + @NoArgsConstructor + class EntryFilterDefinition { + + /** + * The name of the broker interceptor. + */ + private String name; + + /** + * The description of the broker interceptor to be used for user help. + */ + private String description; + + /** + * The class name for the broker interceptor. + */ + private String entryFilterClass; + } + + class EntryFilterWithClassLoader implements EntryFilter { + private final EntryFilter entryFilter; + private final NarClassLoader classLoader; + + public EntryFilterWithClassLoader(EntryFilter entryFilter, NarClassLoader classLoader) { + this.entryFilter = entryFilter; + this.classLoader = classLoader; + } + + @Override + public FilterResult filterEntry(Entry entry, FilterContext context) { + return entryFilter.filterEntry(entry, context); + } + } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java index e258b7fbef48f..a59aa486223de 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java @@ -7,7 +7,7 @@ * "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 + * 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 @@ -19,21 +19,136 @@ package org.apache.pulsar.broker.service; import static com.google.common.base.Preconditions.checkArgument; +import com.google.common.collect.ImmutableMap; +import java.io.File; +import java.io.IOException; +import java.nio.file.DirectoryStream; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.Paths; +import java.util.Collections; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; +import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.intercept.BrokerInterceptor; +import org.apache.pulsar.common.nar.NarClassLoader; +import org.apache.pulsar.common.util.ObjectMapperFactory; +@Slf4j public class EntryFilterProvider { + static final String ENTRY_FILTER_DEFINITION_FILE = "entry_filter.yml"; + /** * create entry filter instance. */ - public static EntryFilter createEntryFilter(String className) { - Class entryFilterClass; + public static ImmutableMap createEntryFilters( + ServiceConfiguration conf) throws IOException { + EntryFilter.EntryFilterDefinitions definitions = searchForEntryFilters(conf.getEntryFiltersDirectory(), + conf.getNarExtractionDirectory()); + ImmutableMap.Builder builder = ImmutableMap.builder(); + conf.getEntryFilterNames().forEach(filterName -> { + EntryFilter.EntryFilterMetaData metaData = definitions.getFilters().get(filterName); + if (null == metaData) { + throw new RuntimeException("No entry filter is found for name `" + filterName + + "`. Available entry filters are : " + definitions.getFilters()); + } + EntryFilter.EntryFilterWithClassLoader filter; + try { + filter = load(metaData, conf.getNarExtractionDirectory()); + if (filter != null) { + builder.put(filterName, filter); + } + log.info("Successfully loaded entry filter for name `{}`", filterName); + } catch (IOException e) { + log.error("Failed to load the entry filter for name `" + filterName + "`", e); + throw new RuntimeException("Failed to load the broker interceptor for name `" + filterName + "`"); + } + }); + return builder.build(); + } + + private static EntryFilter.EntryFilterDefinitions searchForEntryFilters(String entryFiltersDirectory, + String narExtractionDirectory) + throws IOException { + Path path = Paths.get(entryFiltersDirectory).toAbsolutePath(); + log.info("Searching for entry filters in {}", path); + + EntryFilter.EntryFilterDefinitions entryFilterDefinitions = new EntryFilter.EntryFilterDefinitions(); + if (!path.toFile().exists()) { + log.warn("Pulsar entry filters directory not found"); + return entryFilterDefinitions; + } + + try (DirectoryStream stream = Files.newDirectoryStream(path, "*.nar")) { + for (Path archive : stream) { + try { + EntryFilter.EntryFilterDefinition def = + getEntryFilterDefinition(archive.toString(), narExtractionDirectory); + log.info("Found entry filter from {} : {}", archive, def); + + checkArgument(StringUtils.isNotBlank(def.getName())); + checkArgument(StringUtils.isNotBlank(def.getEntryFilterClass())); + + EntryFilter.EntryFilterMetaData metadata = new EntryFilter.EntryFilterMetaData(); + metadata.setDefinition(def); + metadata.setArchivePath(archive); + + entryFilterDefinitions.getFilters().put(def.getName(), metadata); + } catch (Throwable t) { + log.warn("Failed to load entry filters from {}." + + " It is OK however if you want to use this entry filters," + + " please make sure you put the correct entry filter NAR" + + " package in the entry filter directory.", archive, t); + } + } + } + + return entryFilterDefinitions; + } + + private static EntryFilter.EntryFilterDefinition getEntryFilterDefinition(String narPath, + String narExtractionDirectory) + throws IOException { + try (NarClassLoader ncl = NarClassLoader.getFromArchive(new File(narPath), Collections.emptySet(), + narExtractionDirectory)) { + return getEntryFilterDefinition(ncl); + } + } + + private static EntryFilter.EntryFilterDefinition getEntryFilterDefinition(NarClassLoader ncl) throws IOException { + String configStr = ncl.getServiceDefinition(ENTRY_FILTER_DEFINITION_FILE); + + return ObjectMapperFactory.getThreadLocalYaml().readValue( + configStr, EntryFilter.EntryFilterDefinition.class + ); + } + + private static EntryFilter.EntryFilterWithClassLoader load(EntryFilter.EntryFilterMetaData metadata, + String narExtractionDirectory) + throws IOException { + NarClassLoader ncl = NarClassLoader.getFromArchive( + metadata.getArchivePath().toAbsolutePath().toFile(), + Collections.emptySet(), + BrokerInterceptor.class.getClassLoader(), narExtractionDirectory); + + EntryFilter.EntryFilterDefinition def = getEntryFilterDefinition(ncl); + if (StringUtils.isBlank(def.getEntryFilterClass())) { + throw new IOException("Entry filters `" + def.getName() + "` does NOT provide a broker" + + " interceptors implementation"); + } + try { - entryFilterClass = Class.forName(className); - Object obj = entryFilterClass.getDeclaredConstructor().newInstance(); - checkArgument(obj instanceof EntryFilter, "The instance is not an instance of " + className); - return (EntryFilter) obj; - } catch (Exception e) { - throw new RuntimeException(e); + Class entryFilterClass = ncl.loadClass(def.getEntryFilterClass()); + Object filter = entryFilterClass.getDeclaredConstructor().newInstance(); + if (!(filter instanceof EntryFilter)) { + throw new IOException("Class " + def.getEntryFilterClass() + + " does not implement entry filter interface"); + } + EntryFilter pi = (EntryFilter) filter; + return new EntryFilter.EntryFilterWithClassLoader(pi, ncl); + } catch (Throwable t) { + return null; } } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java index 9db4111969b81..89cc4480bf185 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java @@ -294,4 +294,10 @@ default boolean isSystemTopic() { */ CompletableFuture truncate(); + /** + * Get BrokerService. + * @return + */ + BrokerService getBrokerService(); + } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java index c72cabfa73221..8d82ec1fccd4e 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java @@ -20,13 +20,16 @@ import static org.testng.AssertJUnit.assertEquals; import static org.testng.AssertJUnit.assertNotNull; -import java.util.Arrays; +import com.google.common.collect.ImmutableList; +import java.lang.reflect.Field; import java.util.HashMap; import java.util.Map; import java.util.UUID; import java.util.concurrent.TimeUnit; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.pulsar.broker.service.AbstractBaseDispatcher; import org.apache.pulsar.broker.service.BrokerTestBase; +import org.apache.pulsar.broker.service.Dispatcher; import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; @@ -53,15 +56,19 @@ protected void cleanup() throws Exception { } public void testFilter() throws Exception { - internalCleanup(); - conf.setEntryFilterClassNames(Arrays.asList("org.apache.pulsar.broker.service.plugin.EntryFilterForTest", - "org.apache.pulsar.broker.service.plugin.EntryFilterForTest2")); - baseSetup(); String topic = "persistent://prop/ns-abc/topic" + UUID.randomUUID(); String subName = "sub"; Consumer consumer = pulsarClient.newConsumer(Schema.STRING).topic(topic) .subscriptionName(subName).subscribe(); + // mock entry filters + PersistentSubscription subscription = (PersistentSubscription) pulsar.getBrokerService() + .getTopicReference(topic).get().getSubscription(subName); + Dispatcher dispatcher = subscription.getDispatcher(); + Field field = AbstractBaseDispatcher.class.getDeclaredField("entryFilters"); + field.setAccessible(true); + field.set(dispatcher, ImmutableList.of(new EntryFilterForTest(), new EntryFilterForTest2())); + Producer producer = pulsarClient.newProducer(Schema.STRING) .enableBatching(false) .topic(topic).create(); @@ -99,9 +106,6 @@ public void testFilter() throws Exception { assertEquals(0, counter); // All messages should be acked, check the MarkDeletedPosition - PersistentSubscription subscription = - (PersistentSubscription) pulsar.getBrokerService() - .getTopicReference(topic).get().getSubscription(subName); assertNotNull(lastMsgId); MessageIdImpl finalLastMsgId = lastMsgId; Awaitility.await().untilAsserted(() -> { From 02c1e6227dfcbb3405234cd7e47e9185d2bfcfc5 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Fri, 26 Nov 2021 14:38:20 +0800 Subject: [PATCH 10/19] fix check style --- .../org/apache/pulsar/broker/service/EntryFilterProvider.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java index a59aa486223de..eaba9cfd5eb52 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java @@ -7,7 +7,7 @@ * "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 + * 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 From efb7078581749356a1766630152db6e1b757d124 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Fri, 26 Nov 2021 15:17:24 +0800 Subject: [PATCH 11/19] split into different class --- .../service/AbstractBaseDispatcher.java | 11 +- .../pulsar/broker/service/BrokerService.java | 4 +- .../pulsar/broker/service/EntryFilter.java | 132 ------------------ .../broker/service/plugin/EntryFilter.java | 47 +++++++ .../service/plugin/EntryFilterDefinition.java | 42 ++++++ .../plugin/EntryFilterDefinitions.java | 28 ++++ .../service/plugin/EntryFilterMetaData.java | 37 +++++ .../{ => plugin}/EntryFilterProvider.java | 32 ++--- .../plugin/EntryFilterWithClassLoader.java | 19 +++ .../broker/service/plugin/FilterContext.java | 51 +++++++ .../broker/service/plugin/package-info.java | 19 +++ ...terForTest2.java => EntryFilter2Test.java} | 3 +- ...ilterForTest.java => EntryFilterTest.java} | 3 +- .../service/plugin/FilterEntryTest.java | 2 +- 14 files changed, 272 insertions(+), 158 deletions(-) delete mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterDefinition.java create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterDefinitions.java create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterMetaData.java rename pulsar-broker/src/main/java/org/apache/pulsar/broker/service/{ => plugin}/EntryFilterProvider.java (80%) create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterWithClassLoader.java create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/package-info.java rename pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/{EntryFilterForTest2.java => EntryFilter2Test.java} (94%) rename pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/{EntryFilterForTest.java => EntryFilterTest.java} (93%) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index 519612fe33f33..c7f2e70516e81 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -34,6 +34,9 @@ import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.intercept.BrokerInterceptor; import org.apache.pulsar.broker.service.persistent.PersistentTopic; +import org.apache.pulsar.broker.service.plugin.EntryFilter; +import org.apache.pulsar.broker.service.plugin.EntryFilterWithClassLoader; +import org.apache.pulsar.broker.service.plugin.FilterContext; import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.common.api.proto.CommandAck.AckType; import org.apache.pulsar.common.api.proto.MessageMetadata; @@ -55,7 +58,7 @@ public abstract class AbstractBaseDispatcher implements Dispatcher { * Entry filters in Broker. * Not set to final, for the convenience of testing mock. */ - protected ImmutableList entryFilters; + protected ImmutableList entryFilters; protected AbstractBaseDispatcher(Subscription subscription, ServiceConfiguration serviceConfig) { this.subscription = subscription; @@ -126,7 +129,7 @@ public int filterEntriesForConsumer(Optional entryWrapper, int e long totalBytes = 0; int totalChunkedMessages = 0; int totalEntries = 0; - EntryFilter.FilterContext filterContext = new EntryFilter.FilterContext(); + FilterContext filterContext = new FilterContext(); for (int i = 0, entriesSize = entries.size(); i < entriesSize; i++) { Entry entry = entries.get(i); if (entry == null) { @@ -216,7 +219,7 @@ && trackDelayedDelivery(entry.getLedgerId(), entry.getEntryId(), msgMetadata)) { return totalEntries; } - private EntryFilter.FilterResult getFilterResult(EntryFilter.FilterContext filterContext, Entry entry) { + private EntryFilter.FilterResult getFilterResult(FilterContext filterContext, Entry entry) { EntryFilter.FilterResult result = EntryFilter.FilterResult.REJECT; for (EntryFilter entryFilter : entryFilters) { if (entryFilter.filterEntry(entry, filterContext) == EntryFilter.FilterResult.ACCEPT) { @@ -227,7 +230,7 @@ private EntryFilter.FilterResult getFilterResult(EntryFilter.FilterContext filte return result; } - private void fillContext(EntryFilter.FilterContext context, + private void fillContext(FilterContext context, EntryBatchSizes batchSizes, SendMessageInfo sendMessageInfo, EntryBatchIndexesAcks indexesAcks, ManagedCursor cursor, boolean isReplayRead, MessageMetadata msgMetadata, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 2f8f7ab7cc676..6846548b7ddfd 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -115,6 +115,8 @@ import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.broker.service.persistent.SystemTopic; +import org.apache.pulsar.broker.service.plugin.EntryFilterProvider; +import org.apache.pulsar.broker.service.plugin.EntryFilterWithClassLoader; import org.apache.pulsar.broker.stats.ClusterReplicationMetrics; import org.apache.pulsar.broker.stats.prometheus.metrics.ObserverGauge; import org.apache.pulsar.broker.stats.prometheus.metrics.Summary; @@ -268,7 +270,7 @@ public class BrokerService implements Closeable { private boolean preciseTopicPublishRateLimitingEnable; private final LongAdder pausedConnections = new LongAdder(); private BrokerInterceptor interceptor; - private ImmutableMap entryFilters; + private ImmutableMap entryFilters; private Set brokerEntryMetadataInterceptors; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java deleted file mode 100644 index a4f937477b314..0000000000000 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilter.java +++ /dev/null @@ -1,132 +0,0 @@ -/** - * 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.pulsar.broker.service; - -import java.nio.file.Path; -import java.util.Map; -import java.util.TreeMap; -import lombok.Data; -import lombok.NoArgsConstructor; -import org.apache.bookkeeper.mledger.Entry; -import org.apache.bookkeeper.mledger.ManagedCursor; -import org.apache.pulsar.common.api.proto.MessageMetadata; -import org.apache.pulsar.common.nar.NarClassLoader; - -public interface EntryFilter { - - /** - * Broker determines whether to filter out this entry based on the return value of this method. - * Do not deserialize the entire entry in this method, - * which has a great impact on the broker's memory and CPU. - * @param entry - * @param context - * @return - */ - FilterResult filterEntry(Entry entry, FilterContext context); - - @Data - class FilterContext { - - private EntryBatchSizes batchSizes; - private SendMessageInfo sendMessageInfo; - private EntryBatchIndexesAcks indexesAcks; - private ManagedCursor cursor; - private boolean isReplayRead; - private Subscription subscription; - private SubscriptionOption subscriptionOption; - private MessageMetadata msgMetadata; - - public void reset() { - batchSizes = null; - sendMessageInfo = null; - indexesAcks = null; - cursor = null; - isReplayRead = false; - subscription = null; - subscriptionOption = null; - msgMetadata = null; - } - } - - enum FilterResult { - /** - * deliver to the consumer. - */ - ACCEPT, - /** - * skip the message. - */ - REJECT, - } - - @Data - class EntryFilterDefinitions { - private final Map filters = new TreeMap<>(); - } - - @Data - @NoArgsConstructor - class EntryFilterMetaData { - - /** - * The definition of the broker interceptor. - */ - private EntryFilterDefinition definition; - - /** - * The path to the handler package. - */ - private Path archivePath; - } - - @Data - @NoArgsConstructor - class EntryFilterDefinition { - - /** - * The name of the broker interceptor. - */ - private String name; - - /** - * The description of the broker interceptor to be used for user help. - */ - private String description; - - /** - * The class name for the broker interceptor. - */ - private String entryFilterClass; - } - - class EntryFilterWithClassLoader implements EntryFilter { - private final EntryFilter entryFilter; - private final NarClassLoader classLoader; - - public EntryFilterWithClassLoader(EntryFilter entryFilter, NarClassLoader classLoader) { - this.entryFilter = entryFilter; - this.classLoader = classLoader; - } - - @Override - public FilterResult filterEntry(Entry entry, FilterContext context) { - return entryFilter.filterEntry(entry, context); - } - } -} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java new file mode 100644 index 0000000000000..fd8a4f29e88f2 --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java @@ -0,0 +1,47 @@ +/** + * 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.pulsar.broker.service.plugin; + +import org.apache.bookkeeper.mledger.Entry; + +public interface EntryFilter { + + /** + * Broker determines whether to filter out this entry based on the return value of this method. + * Do not deserialize the entire entry in this method, + * which has a great impact on the broker's memory and CPU. + * @param entry + * @param context + * @return + */ + FilterResult filterEntry(Entry entry, FilterContext context); + + + enum FilterResult { + /** + * deliver to the consumer. + */ + ACCEPT, + /** + * skip the message. + */ + REJECT, + } + +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterDefinition.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterDefinition.java new file mode 100644 index 0000000000000..2f4399a90ac88 --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterDefinition.java @@ -0,0 +1,42 @@ +/** + * 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.pulsar.broker.service.plugin; + +import lombok.Data; +import lombok.NoArgsConstructor; + +@Data +@NoArgsConstructor +public class EntryFilterDefinition { + + /** + * The name of the broker interceptor. + */ + private String name; + + /** + * The description of the broker interceptor to be used for user help. + */ + private String description; + + /** + * The class name for the broker interceptor. + */ + private String entryFilterClass; +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterDefinitions.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterDefinitions.java new file mode 100644 index 0000000000000..9aa3113e177ae --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterDefinitions.java @@ -0,0 +1,28 @@ +/** + * 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.pulsar.broker.service.plugin; + +import java.util.Map; +import java.util.TreeMap; +import lombok.Data; + +@Data +public class EntryFilterDefinitions { + private final Map filters = new TreeMap<>(); +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterMetaData.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterMetaData.java new file mode 100644 index 0000000000000..1756e0bae5ba8 --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterMetaData.java @@ -0,0 +1,37 @@ +/** + * 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.pulsar.broker.service.plugin; + +import java.nio.file.Path; +import lombok.Data; +import lombok.NoArgsConstructor; + +@Data +@NoArgsConstructor +public class EntryFilterMetaData { + /** + * The definition of the broker interceptor. + */ + private EntryFilterDefinition definition; + + /** + * The path to the handler package. + */ + private Path archivePath; +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java similarity index 80% rename from pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java rename to pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java index eaba9cfd5eb52..355f91d326481 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/EntryFilterProvider.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java @@ -16,7 +16,7 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.pulsar.broker.service; +package org.apache.pulsar.broker.service.plugin; import static com.google.common.base.Preconditions.checkArgument; import com.google.common.collect.ImmutableMap; @@ -42,18 +42,18 @@ public class EntryFilterProvider { /** * create entry filter instance. */ - public static ImmutableMap createEntryFilters( + public static ImmutableMap createEntryFilters( ServiceConfiguration conf) throws IOException { - EntryFilter.EntryFilterDefinitions definitions = searchForEntryFilters(conf.getEntryFiltersDirectory(), + EntryFilterDefinitions definitions = searchForEntryFilters(conf.getEntryFiltersDirectory(), conf.getNarExtractionDirectory()); - ImmutableMap.Builder builder = ImmutableMap.builder(); + ImmutableMap.Builder builder = ImmutableMap.builder(); conf.getEntryFilterNames().forEach(filterName -> { - EntryFilter.EntryFilterMetaData metaData = definitions.getFilters().get(filterName); + EntryFilterMetaData metaData = definitions.getFilters().get(filterName); if (null == metaData) { throw new RuntimeException("No entry filter is found for name `" + filterName + "`. Available entry filters are : " + definitions.getFilters()); } - EntryFilter.EntryFilterWithClassLoader filter; + EntryFilterWithClassLoader filter; try { filter = load(metaData, conf.getNarExtractionDirectory()); if (filter != null) { @@ -68,13 +68,13 @@ public static ImmutableMap creat return builder.build(); } - private static EntryFilter.EntryFilterDefinitions searchForEntryFilters(String entryFiltersDirectory, + private static EntryFilterDefinitions searchForEntryFilters(String entryFiltersDirectory, String narExtractionDirectory) throws IOException { Path path = Paths.get(entryFiltersDirectory).toAbsolutePath(); log.info("Searching for entry filters in {}", path); - EntryFilter.EntryFilterDefinitions entryFilterDefinitions = new EntryFilter.EntryFilterDefinitions(); + EntryFilterDefinitions entryFilterDefinitions = new EntryFilterDefinitions(); if (!path.toFile().exists()) { log.warn("Pulsar entry filters directory not found"); return entryFilterDefinitions; @@ -83,14 +83,14 @@ private static EntryFilter.EntryFilterDefinitions searchForEntryFilters(String e try (DirectoryStream stream = Files.newDirectoryStream(path, "*.nar")) { for (Path archive : stream) { try { - EntryFilter.EntryFilterDefinition def = + EntryFilterDefinition def = getEntryFilterDefinition(archive.toString(), narExtractionDirectory); log.info("Found entry filter from {} : {}", archive, def); checkArgument(StringUtils.isNotBlank(def.getName())); checkArgument(StringUtils.isNotBlank(def.getEntryFilterClass())); - EntryFilter.EntryFilterMetaData metadata = new EntryFilter.EntryFilterMetaData(); + EntryFilterMetaData metadata = new EntryFilterMetaData(); metadata.setDefinition(def); metadata.setArchivePath(archive); @@ -107,7 +107,7 @@ private static EntryFilter.EntryFilterDefinitions searchForEntryFilters(String e return entryFilterDefinitions; } - private static EntryFilter.EntryFilterDefinition getEntryFilterDefinition(String narPath, + private static EntryFilterDefinition getEntryFilterDefinition(String narPath, String narExtractionDirectory) throws IOException { try (NarClassLoader ncl = NarClassLoader.getFromArchive(new File(narPath), Collections.emptySet(), @@ -116,15 +116,15 @@ private static EntryFilter.EntryFilterDefinition getEntryFilterDefinition(String } } - private static EntryFilter.EntryFilterDefinition getEntryFilterDefinition(NarClassLoader ncl) throws IOException { + private static EntryFilterDefinition getEntryFilterDefinition(NarClassLoader ncl) throws IOException { String configStr = ncl.getServiceDefinition(ENTRY_FILTER_DEFINITION_FILE); return ObjectMapperFactory.getThreadLocalYaml().readValue( - configStr, EntryFilter.EntryFilterDefinition.class + configStr, EntryFilterDefinition.class ); } - private static EntryFilter.EntryFilterWithClassLoader load(EntryFilter.EntryFilterMetaData metadata, + private static EntryFilterWithClassLoader load(EntryFilterMetaData metadata, String narExtractionDirectory) throws IOException { NarClassLoader ncl = NarClassLoader.getFromArchive( @@ -132,7 +132,7 @@ private static EntryFilter.EntryFilterWithClassLoader load(EntryFilter.EntryFilt Collections.emptySet(), BrokerInterceptor.class.getClassLoader(), narExtractionDirectory); - EntryFilter.EntryFilterDefinition def = getEntryFilterDefinition(ncl); + EntryFilterDefinition def = getEntryFilterDefinition(ncl); if (StringUtils.isBlank(def.getEntryFilterClass())) { throw new IOException("Entry filters `" + def.getName() + "` does NOT provide a broker" + " interceptors implementation"); @@ -146,7 +146,7 @@ private static EntryFilter.EntryFilterWithClassLoader load(EntryFilter.EntryFilt + " does not implement entry filter interface"); } EntryFilter pi = (EntryFilter) filter; - return new EntryFilter.EntryFilterWithClassLoader(pi, ncl); + return new EntryFilterWithClassLoader(pi, ncl); } catch (Throwable t) { return null; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterWithClassLoader.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterWithClassLoader.java new file mode 100644 index 0000000000000..7f0b648e9e240 --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterWithClassLoader.java @@ -0,0 +1,19 @@ +package org.apache.pulsar.broker.service.plugin; + +import org.apache.bookkeeper.mledger.Entry; +import org.apache.pulsar.common.nar.NarClassLoader; + +public class EntryFilterWithClassLoader implements EntryFilter { + private final EntryFilter entryFilter; + private final NarClassLoader classLoader; + + public EntryFilterWithClassLoader(EntryFilter entryFilter, NarClassLoader classLoader) { + this.entryFilter = entryFilter; + this.classLoader = classLoader; + } + + @Override + public FilterResult filterEntry(Entry entry, FilterContext context) { + return entryFilter.filterEntry(entry, context); + } +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java new file mode 100644 index 0000000000000..54353d8660204 --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java @@ -0,0 +1,51 @@ +/** + * 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.pulsar.broker.service.plugin; + +import lombok.Data; +import org.apache.bookkeeper.mledger.ManagedCursor; +import org.apache.pulsar.broker.service.EntryBatchIndexesAcks; +import org.apache.pulsar.broker.service.EntryBatchSizes; +import org.apache.pulsar.broker.service.SendMessageInfo; +import org.apache.pulsar.broker.service.Subscription; +import org.apache.pulsar.broker.service.SubscriptionOption; +import org.apache.pulsar.common.api.proto.MessageMetadata; + +@Data +public class FilterContext { + private EntryBatchSizes batchSizes; + private SendMessageInfo sendMessageInfo; + private EntryBatchIndexesAcks indexesAcks; + private ManagedCursor cursor; + private boolean isReplayRead; + private Subscription subscription; + private SubscriptionOption subscriptionOption; + private MessageMetadata msgMetadata; + + public void reset() { + batchSizes = null; + sendMessageInfo = null; + indexesAcks = null; + cursor = null; + isReplayRead = false; + subscription = null; + subscriptionOption = null; + msgMetadata = null; + } +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/package-info.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/package-info.java new file mode 100644 index 0000000000000..05e54ca4831f8 --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/package-info.java @@ -0,0 +1,19 @@ +/** + * 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.pulsar.broker.service.plugin; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest2.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilter2Test.java similarity index 94% rename from pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest2.java rename to pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilter2Test.java index b6ab2afb5fe5b..c6e3f4d27cafa 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest2.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilter2Test.java @@ -22,11 +22,10 @@ import java.util.List; import org.apache.bookkeeper.mledger.Entry; import org.apache.commons.collections4.MapUtils; -import org.apache.pulsar.broker.service.EntryFilter; import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.common.api.proto.KeyValue; -public class EntryFilterForTest2 implements EntryFilter { +public class EntryFilter2Test implements EntryFilter { @Override public FilterResult filterEntry(Entry entry, FilterContext context) { if (context.getMsgMetadata() == null || context.getMsgMetadata().getPropertiesCount() <= 0) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterTest.java similarity index 93% rename from pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java rename to pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterTest.java index e5640e5c7f3c8..ccef93799b5a6 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterForTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterTest.java @@ -21,10 +21,9 @@ import java.util.List; import org.apache.bookkeeper.mledger.Entry; -import org.apache.pulsar.broker.service.EntryFilter; import org.apache.pulsar.common.api.proto.KeyValue; -public class EntryFilterForTest implements EntryFilter { +public class EntryFilterTest implements EntryFilter { @Override public FilterResult filterEntry(Entry entry, FilterContext context) { if (context.getMsgMetadata() == null || context.getMsgMetadata().getPropertiesCount() <= 0) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java index 8d82ec1fccd4e..5361734edffb9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java @@ -67,7 +67,7 @@ public void testFilter() throws Exception { Dispatcher dispatcher = subscription.getDispatcher(); Field field = AbstractBaseDispatcher.class.getDeclaredField("entryFilters"); field.setAccessible(true); - field.set(dispatcher, ImmutableList.of(new EntryFilterForTest(), new EntryFilterForTest2())); + field.set(dispatcher, ImmutableList.of(new EntryFilterTest(), new EntryFilter2Test())); Producer producer = pulsarClient.newProducer(Schema.STRING) .enableBatching(false) From 313d4023dc37e80879925752e6da90f384c31541 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Fri, 26 Nov 2021 15:26:52 +0800 Subject: [PATCH 12/19] fix check style --- .../plugin/EntryFilterWithClassLoader.java | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterWithClassLoader.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterWithClassLoader.java index 7f0b648e9e240..c43620e1a692f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterWithClassLoader.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterWithClassLoader.java @@ -1,3 +1,21 @@ +/** + * 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.pulsar.broker.service.plugin; import org.apache.bookkeeper.mledger.Entry; From 0d9f1a79ecbd1aac7b2c75aa35567b4510302376 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Fri, 26 Nov 2021 15:52:26 +0800 Subject: [PATCH 13/19] fix unit test --- .../apache/pulsar/broker/service/AbstractBaseDispatcher.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index c7f2e70516e81..6fd999877b651 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -64,7 +64,8 @@ protected AbstractBaseDispatcher(Subscription subscription, ServiceConfiguration this.subscription = subscription; this.serviceConfig = serviceConfig; this.dispatchThrottlingOnBatchMessageEnabled = serviceConfig.isDispatchThrottlingOnBatchMessageEnabled(); - if (MapUtils.isNotEmpty(subscription.getTopic().getBrokerService().getEntryFilters())) { + if (subscription != null && MapUtils.isNotEmpty(subscription.getTopic() + .getBrokerService().getEntryFilters())) { this.entryFilters = subscription.getTopic().getBrokerService().getEntryFilters().values().asList(); } else { this.entryFilters = ImmutableList.of(); From b5611dd04b35ff78af9d1292e92ce334b7e643b3 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Fri, 26 Nov 2021 16:35:43 +0800 Subject: [PATCH 14/19] fix unit test --- .../apache/pulsar/broker/service/AbstractBaseDispatcher.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index 6fd999877b651..84471f19243a8 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -64,7 +64,7 @@ protected AbstractBaseDispatcher(Subscription subscription, ServiceConfiguration this.subscription = subscription; this.serviceConfig = serviceConfig; this.dispatchThrottlingOnBatchMessageEnabled = serviceConfig.isDispatchThrottlingOnBatchMessageEnabled(); - if (subscription != null && MapUtils.isNotEmpty(subscription.getTopic() + if (subscription != null && subscription.getTopic() != null && MapUtils.isNotEmpty(subscription.getTopic() .getBrokerService().getEntryFilters())) { this.entryFilters = subscription.getTopic().getBrokerService().getEntryFilters().values().asList(); } else { From b428669e3e9fb12e404b9c20bb52f78fe7ca6a35 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Fri, 26 Nov 2021 17:33:59 +0800 Subject: [PATCH 15/19] address comments --- .../service/AbstractBaseDispatcher.java | 38 +++++++++---------- .../service/plugin/EntryFilterProvider.java | 17 +++------ .../broker/service/plugin/FilterContext.java | 14 +------ 3 files changed, 26 insertions(+), 43 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index 84471f19243a8..1fb8c6bfb352a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -21,12 +21,14 @@ import com.google.common.collect.ImmutableList; import io.netty.buffer.ByteBuf; +import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Optional; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; +import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.commons.collections4.CollectionUtils; import org.apache.commons.collections4.MapUtils; @@ -59,6 +61,7 @@ public abstract class AbstractBaseDispatcher implements Dispatcher { * Not set to final, for the convenience of testing mock. */ protected ImmutableList entryFilters; + protected final FilterContext filterContext; protected AbstractBaseDispatcher(Subscription subscription, ServiceConfiguration serviceConfig) { this.subscription = subscription; @@ -67,8 +70,10 @@ protected AbstractBaseDispatcher(Subscription subscription, ServiceConfiguration if (subscription != null && subscription.getTopic() != null && MapUtils.isNotEmpty(subscription.getTopic() .getBrokerService().getEntryFilters())) { this.entryFilters = subscription.getTopic().getBrokerService().getEntryFilters().values().asList(); + this.filterContext = new FilterContext(); } else { this.entryFilters = ImmutableList.of(); + this.filterContext = FilterContext.FILTER_CONTEXT_DISABLED; } } @@ -130,7 +135,7 @@ public int filterEntriesForConsumer(Optional entryWrapper, int e long totalBytes = 0; int totalChunkedMessages = 0; int totalEntries = 0; - FilterContext filterContext = new FilterContext(); + List entriesToFiltered = CollectionUtils.isNotEmpty(entryFilters) ? new ArrayList<>() : null; for (int i = 0, entriesSize = entries.size(); i < entriesSize; i++) { Entry entry = entries.get(i); if (entry == null) { @@ -146,15 +151,11 @@ public int filterEntriesForConsumer(Optional entryWrapper, int e ? Commands.peekMessageMetadata(metadataAndPayload, subscription.toString(), -1) : msgMetadata; if (CollectionUtils.isNotEmpty(entryFilters)) { - fillContext(filterContext, batchSizes, sendMessageInfo, indexesAcks, cursor, isReplayRead, - msgMetadata, subscription); - EntryFilter.FilterResult result = getFilterResult(filterContext, entry); - if (EntryFilter.FilterResult.REJECT == result) { - PositionImpl pos = (PositionImpl) entry.getPosition(); + fillContext(filterContext, msgMetadata, subscription); + if (EntryFilter.FilterResult.REJECT == getFilterResult(filterContext, entry)) { + entriesToFiltered.add(entry.getPosition()); entries.set(i, null); entry.release(); - subscription.acknowledgeMessage(Collections.singletonList(pos), AckType.Individual, - Collections.emptyMap()); continue; } } @@ -214,6 +215,11 @@ && trackDelayedDelivery(entry.getLedgerId(), entry.getEntryId(), msgMetadata)) { interceptor.beforeSendMessage(subscription, entry, ackSet, msgMetadata); } } + if (CollectionUtils.isNotEmpty(entriesToFiltered)) { + subscription.acknowledgeMessage(entriesToFiltered, AckType.Individual, + Collections.emptyMap()); + } + sendMessageInfo.setTotalMessages(totalMessages); sendMessageInfo.setTotalBytes(totalBytes); sendMessageInfo.setTotalChunkedMessages(totalChunkedMessages); @@ -221,27 +227,19 @@ && trackDelayedDelivery(entry.getLedgerId(), entry.getEntryId(), msgMetadata)) { } private EntryFilter.FilterResult getFilterResult(FilterContext filterContext, Entry entry) { - EntryFilter.FilterResult result = EntryFilter.FilterResult.REJECT; + EntryFilter.FilterResult result = EntryFilter.FilterResult.ACCEPT; for (EntryFilter entryFilter : entryFilters) { - if (entryFilter.filterEntry(entry, filterContext) == EntryFilter.FilterResult.ACCEPT) { - result = EntryFilter.FilterResult.ACCEPT; + if (entryFilter.filterEntry(entry, filterContext) == EntryFilter.FilterResult.REJECT) { + result = EntryFilter.FilterResult.REJECT; break; } } return result; } - private void fillContext(FilterContext context, - EntryBatchSizes batchSizes, SendMessageInfo sendMessageInfo, - EntryBatchIndexesAcks indexesAcks, ManagedCursor cursor, - boolean isReplayRead, MessageMetadata msgMetadata, + private void fillContext(FilterContext context, MessageMetadata msgMetadata, Subscription subscription) { context.reset(); - context.setBatchSizes(batchSizes); - context.setSendMessageInfo(sendMessageInfo); - context.setIndexesAcks(indexesAcks); - context.setCursor(cursor); - context.setReplayRead(isReplayRead); context.setMsgMetadata(msgMetadata); context.setSubscription(subscription); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java index 355f91d326481..3de5334477793 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java @@ -47,24 +47,19 @@ public static ImmutableMap createEntryFilter EntryFilterDefinitions definitions = searchForEntryFilters(conf.getEntryFiltersDirectory(), conf.getNarExtractionDirectory()); ImmutableMap.Builder builder = ImmutableMap.builder(); - conf.getEntryFilterNames().forEach(filterName -> { + for (String filterName : conf.getEntryFilterNames()) { EntryFilterMetaData metaData = definitions.getFilters().get(filterName); if (null == metaData) { throw new RuntimeException("No entry filter is found for name `" + filterName + "`. Available entry filters are : " + definitions.getFilters()); } EntryFilterWithClassLoader filter; - try { - filter = load(metaData, conf.getNarExtractionDirectory()); - if (filter != null) { - builder.put(filterName, filter); - } - log.info("Successfully loaded entry filter for name `{}`", filterName); - } catch (IOException e) { - log.error("Failed to load the entry filter for name `" + filterName + "`", e); - throw new RuntimeException("Failed to load the broker interceptor for name `" + filterName + "`"); + filter = load(metaData, conf.getNarExtractionDirectory()); + if (filter != null) { + builder.put(filterName, filter); } - }); + log.info("Successfully loaded entry filter for name `{}`", filterName); + } return builder.build(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java index 54353d8660204..b392c7e0d0655 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java @@ -29,23 +29,13 @@ @Data public class FilterContext { - private EntryBatchSizes batchSizes; - private SendMessageInfo sendMessageInfo; - private EntryBatchIndexesAcks indexesAcks; - private ManagedCursor cursor; - private boolean isReplayRead; private Subscription subscription; - private SubscriptionOption subscriptionOption; private MessageMetadata msgMetadata; public void reset() { - batchSizes = null; - sendMessageInfo = null; - indexesAcks = null; - cursor = null; - isReplayRead = false; subscription = null; - subscriptionOption = null; msgMetadata = null; } + + public static FilterContext FILTER_CONTEXT_DISABLED = new FilterContext(); } From d69a0fda077938b4c0c5066fc9c487a828e11671 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Fri, 26 Nov 2021 17:37:07 +0800 Subject: [PATCH 16/19] change doc --- .../pulsar/broker/service/plugin/EntryFilterDefinition.java | 6 +++--- .../pulsar/broker/service/plugin/EntryFilterMetaData.java | 2 +- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterDefinition.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterDefinition.java index 2f4399a90ac88..5df3944e41dd2 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterDefinition.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterDefinition.java @@ -26,17 +26,17 @@ public class EntryFilterDefinition { /** - * The name of the broker interceptor. + * The name of the entry filter. */ private String name; /** - * The description of the broker interceptor to be used for user help. + * The description of the entry filter to be used for user help. */ private String description; /** - * The class name for the broker interceptor. + * The class name for the entry filter. */ private String entryFilterClass; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterMetaData.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterMetaData.java index 1756e0bae5ba8..babaa80d5f60e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterMetaData.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterMetaData.java @@ -26,7 +26,7 @@ @NoArgsConstructor public class EntryFilterMetaData { /** - * The definition of the broker interceptor. + * The definition of the entry filter. */ private EntryFilterDefinition definition; From e7c51a48e845ea8569a00044f747ee3cdfe6c2ad Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Fri, 26 Nov 2021 17:55:44 +0800 Subject: [PATCH 17/19] fix check style --- .../apache/pulsar/broker/service/plugin/FilterContext.java | 7 +------ 1 file changed, 1 insertion(+), 6 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java index b392c7e0d0655..e520e1011f9ff 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/FilterContext.java @@ -19,12 +19,7 @@ package org.apache.pulsar.broker.service.plugin; import lombok.Data; -import org.apache.bookkeeper.mledger.ManagedCursor; -import org.apache.pulsar.broker.service.EntryBatchIndexesAcks; -import org.apache.pulsar.broker.service.EntryBatchSizes; -import org.apache.pulsar.broker.service.SendMessageInfo; import org.apache.pulsar.broker.service.Subscription; -import org.apache.pulsar.broker.service.SubscriptionOption; import org.apache.pulsar.common.api.proto.MessageMetadata; @Data @@ -37,5 +32,5 @@ public void reset() { msgMetadata = null; } - public static FilterContext FILTER_CONTEXT_DISABLED = new FilterContext(); + public static final FilterContext FILTER_CONTEXT_DISABLED = new FilterContext(); } From e16c16b88355d4f5cf0256170bd0f4ff96ac5962 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Sat, 27 Nov 2021 09:14:45 +0800 Subject: [PATCH 18/19] add close --- .../pulsar/broker/service/plugin/EntryFilter.java | 5 +++++ .../service/plugin/EntryFilterProvider.java | 15 +++++++++------ .../plugin/EntryFilterWithClassLoader.java | 13 +++++++++++++ .../broker/service/plugin/EntryFilter2Test.java | 5 +++++ .../broker/service/plugin/EntryFilterTest.java | 5 +++++ 5 files changed, 37 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java index fd8a4f29e88f2..dac20dd2a5079 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java @@ -32,6 +32,11 @@ public interface EntryFilter { */ FilterResult filterEntry(Entry entry, FilterContext context); + /** + * close the entry filter. + */ + void close(); + enum FilterResult { /** diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java index 3de5334477793..b6e713820f5d0 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java @@ -30,7 +30,6 @@ import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.broker.ServiceConfiguration; -import org.apache.pulsar.broker.intercept.BrokerInterceptor; import org.apache.pulsar.common.nar.NarClassLoader; import org.apache.pulsar.common.util.ObjectMapperFactory; @@ -125,12 +124,12 @@ private static EntryFilterWithClassLoader load(EntryFilterMetaData metadata, NarClassLoader ncl = NarClassLoader.getFromArchive( metadata.getArchivePath().toAbsolutePath().toFile(), Collections.emptySet(), - BrokerInterceptor.class.getClassLoader(), narExtractionDirectory); + EntryFilter.class.getClassLoader(), narExtractionDirectory); EntryFilterDefinition def = getEntryFilterDefinition(ncl); if (StringUtils.isBlank(def.getEntryFilterClass())) { - throw new IOException("Entry filters `" + def.getName() + "` does NOT provide a broker" - + " interceptors implementation"); + throw new IOException("Entry filters `" + def.getName() + "` does NOT provide a entry" + + " filters implementation"); } try { @@ -142,8 +141,12 @@ private static EntryFilterWithClassLoader load(EntryFilterMetaData metadata, } EntryFilter pi = (EntryFilter) filter; return new EntryFilterWithClassLoader(pi, ncl); - } catch (Throwable t) { - return null; + } catch (Exception e) { + if (e instanceof IOException) { + throw (IOException) e; + } + log.error("Failed to load class {}", def.getEntryFilterClass(), e); + throw new IOException(e); } } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterWithClassLoader.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterWithClassLoader.java index c43620e1a692f..8c2569ca6335b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterWithClassLoader.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterWithClassLoader.java @@ -18,9 +18,12 @@ */ package org.apache.pulsar.broker.service.plugin; +import java.io.IOException; +import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.Entry; import org.apache.pulsar.common.nar.NarClassLoader; +@Slf4j public class EntryFilterWithClassLoader implements EntryFilter { private final EntryFilter entryFilter; private final NarClassLoader classLoader; @@ -34,4 +37,14 @@ public EntryFilterWithClassLoader(EntryFilter entryFilter, NarClassLoader classL public FilterResult filterEntry(Entry entry, FilterContext context) { return entryFilter.filterEntry(entry, context); } + + @Override + public void close() { + entryFilter.close(); + try { + classLoader.close(); + } catch (IOException e) { + log.error("close EntryFilterWithClassLoader failed", e); + } + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilter2Test.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilter2Test.java index c6e3f4d27cafa..dbcaa160b2391 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilter2Test.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilter2Test.java @@ -43,4 +43,9 @@ public FilterResult filterEntry(Entry entry, FilterContext context) { } return FilterResult.REJECT; } + + @Override + public void close() { + + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterTest.java index ccef93799b5a6..812d49aa7b44e 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/EntryFilterTest.java @@ -40,4 +40,9 @@ public FilterResult filterEntry(Entry entry, FilterContext context) { } return null; } + + @Override + public void close() { + + } } From 74aea663038ecb1f535f7266b3f7898f1ea94336 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Tue, 30 Nov 2021 17:16:07 +0800 Subject: [PATCH 19/19] Add unit test for closing --- .../service/AbstractBaseDispatcher.java | 5 ++-- .../pulsar/broker/service/BrokerService.java | 11 ++++++++ .../broker/service/plugin/EntryFilter.java | 5 ++-- .../service/plugin/EntryFilterProvider.java | 2 +- .../service/plugin/FilterEntryTest.java | 26 ++++++++++++++++--- 5 files changed, 40 insertions(+), 9 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index 1fb8c6bfb352a..8970ec6359028 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -152,7 +152,7 @@ public int filterEntriesForConsumer(Optional entryWrapper, int e : msgMetadata; if (CollectionUtils.isNotEmpty(entryFilters)) { fillContext(filterContext, msgMetadata, subscription); - if (EntryFilter.FilterResult.REJECT == getFilterResult(filterContext, entry)) { + if (EntryFilter.FilterResult.REJECT == getFilterResult(filterContext, entry, entryFilters)) { entriesToFiltered.add(entry.getPosition()); entries.set(i, null); entry.release(); @@ -226,7 +226,8 @@ && trackDelayedDelivery(entry.getLedgerId(), entry.getEntryId(), msgMetadata)) { return totalEntries; } - private EntryFilter.FilterResult getFilterResult(FilterContext filterContext, Entry entry) { + private static EntryFilter.FilterResult getFilterResult(FilterContext filterContext, Entry entry, + ImmutableList entryFilters) { EntryFilter.FilterResult result = EntryFilter.FilterResult.ACCEPT; for (EntryFilter entryFilter : entryFilters) { if (entryFilter.filterEntry(entry, filterContext) == EntryFilter.FilterResult.REJECT) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 6846548b7ddfd..affb85e9137bf 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -710,6 +710,17 @@ public CompletableFuture closeAsync() { } }); + //close entry filters + if (entryFilters != null) { + entryFilters.forEach((name, filter) -> { + try { + filter.close(); + } catch (Exception e) { + log.warn("Error shutting down entry filter {}", name, e); + } + }); + } + CompletableFuture> cancellableDownstreamFutureReference = new CompletableFuture<>(); log.info("Event loops shutting down gracefully..."); List> shutdownEventLoops = new ArrayList<>(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java index dac20dd2a5079..40e6644953f52 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilter.java @@ -23,9 +23,10 @@ public interface EntryFilter { /** - * Broker determines whether to filter out this entry based on the return value of this method. - * Do not deserialize the entire entry in this method, + * 1. Broker determines whether to filter out this entry based on the return value of this method. + * 2. Do not deserialize the entire entry in this method, * which has a great impact on the broker's memory and CPU. + * 3. Return ACCEPT or null will be regarded as ACCEPT. * @param entry * @param context * @return diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java index b6e713820f5d0..9e19a5784d3fc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/plugin/EntryFilterProvider.java @@ -70,7 +70,7 @@ private static EntryFilterDefinitions searchForEntryFilters(String entryFiltersD EntryFilterDefinitions entryFilterDefinitions = new EntryFilterDefinitions(); if (!path.toFile().exists()) { - log.warn("Pulsar entry filters directory not found"); + log.info("Pulsar entry filters directory not found"); return entryFilterDefinitions; } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java index 5361734edffb9..03389352629e6 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/plugin/FilterEntryTest.java @@ -18,9 +18,14 @@ */ package org.apache.pulsar.broker.service.plugin; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import static org.testng.AssertJUnit.assertEquals; import static org.testng.AssertJUnit.assertNotNull; import com.google.common.collect.ImmutableList; +import com.google.common.collect.ImmutableMap; import java.lang.reflect.Field; import java.util.HashMap; import java.util.Map; @@ -28,6 +33,7 @@ import java.util.concurrent.TimeUnit; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.pulsar.broker.service.AbstractBaseDispatcher; +import org.apache.pulsar.broker.service.BrokerService; import org.apache.pulsar.broker.service.BrokerTestBase; import org.apache.pulsar.broker.service.Dispatcher; import org.apache.pulsar.broker.service.persistent.PersistentSubscription; @@ -36,6 +42,7 @@ import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.impl.MessageIdImpl; +import org.apache.pulsar.common.nar.NarClassLoader; import org.awaitility.Awaitility; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -67,11 +74,15 @@ public void testFilter() throws Exception { Dispatcher dispatcher = subscription.getDispatcher(); Field field = AbstractBaseDispatcher.class.getDeclaredField("entryFilters"); field.setAccessible(true); - field.set(dispatcher, ImmutableList.of(new EntryFilterTest(), new EntryFilter2Test())); + NarClassLoader narClassLoader = mock(NarClassLoader.class); + EntryFilter filter1 = new EntryFilterTest(); + EntryFilterWithClassLoader loader1 = spy(new EntryFilterWithClassLoader(filter1, narClassLoader)); + EntryFilter filter2 = new EntryFilter2Test(); + EntryFilterWithClassLoader loader2 = spy(new EntryFilterWithClassLoader(filter2, narClassLoader)); + field.set(dispatcher, ImmutableList.of(loader1, loader2)); Producer producer = pulsarClient.newProducer(Schema.STRING) - .enableBatching(false) - .topic(topic).create(); + .enableBatching(false).topic(topic).create(); for (int i = 0; i < 10; i++) { producer.send("test"); } @@ -138,7 +149,14 @@ public void testFilter() throws Exception { producer.close(); consumer.close(); - } + BrokerService brokerService = pulsar.getBrokerService(); + Field field1 = BrokerService.class.getDeclaredField("entryFilters"); + field1.setAccessible(true); + field1.set(brokerService, ImmutableMap.of("1", loader1, "2", loader2)); + cleanup(); + verify(loader1, times(1)).close(); + verify(loader2, times(1)).close(); + } }