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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions conf/broker.conf
Original file line number Diff line number Diff line change
Expand Up @@ -685,6 +685,9 @@ bookkeeperDiskWeightBasedPlacementEnabled=false
# A value of '0' disables sending any explicit LACs. Default is 0.
bookkeeperExplicitLacIntervalInMills=0

# Expose bookkeeper client managed ledger stats to prometheus. default is false
# bookkeeperClientExposeStatsToPrometheus=false

### --- Managed Ledger --- ###

# Number of bookies to use when creating a ledger
Expand Down Expand Up @@ -787,6 +790,12 @@ managedLedgerReadEntryTimeoutSeconds=0
# Add entry timeout when broker tries to publish message to bookkeeper (0 to disable it).
managedLedgerAddEntryTimeoutSeconds=0

# Managed ledger prometheus stats latency rollover seconds (default: 60s)
managedLedgerPrometheusStatsLatencyRolloverSeconds=60

# Whether trace managed ledger task execution time
managedLedgerTraceTaskExecution=true

# New entries check delay for the cursor under the managed ledger.
# If no new messages in the topic, the cursor will try to check again after the delay time.
# For consumption latency sensitive scenario, can set to a smaller value or set to 0.
Expand Down
9 changes: 9 additions & 0 deletions conf/standalone.conf
Original file line number Diff line number Diff line change
Expand Up @@ -473,6 +473,9 @@ bookkeeperDiskWeightBasedPlacementEnabled=false
# A value of '0' disables sending any explicit LACs. Default is 0.
bookkeeperExplicitLacIntervalInMills=0

# Expose bookkeeper client managed ledger stats to prometheus. default is false
# bookkeeperClientExposeStatsToPrometheus=false

### --- Managed Ledger --- ###

# Number of bookies to use when creating a ledger
Expand Down Expand Up @@ -576,6 +579,12 @@ managedLedgerNewEntriesCheckDelayInMillis=10
# Use Open Range-Set to cache unacked messages
managedLedgerUnackedRangesOpenCacheSetEnabled=true

# Managed ledger prometheus stats latency rollover seconds (default: 60s)
managedLedgerPrometheusStatsLatencyRolloverSeconds=60

# Whether trace managed ledger task execution time
managedLedgerTraceTaskExecution=true

### --- Load balancer --- ###

loadManagerClassName=org.apache.pulsar.broker.loadbalance.NoopLoadManager
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import org.apache.bookkeeper.mledger.AsyncCallbacks.ManagedLedgerInfoCallback;
import org.apache.bookkeeper.mledger.AsyncCallbacks.OpenLedgerCallback;
import org.apache.bookkeeper.mledger.AsyncCallbacks.OpenReadOnlyCursorCallback;
import org.apache.bookkeeper.stats.StatsProvider;

/**
* A factory to open/create managed ledgers and delete them.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,4 +51,19 @@ public class ManagedLedgerFactoryConfig {
* Whether we should make a copy of the entry payloads when inserting in cache
*/
private boolean copyEntriesInCache = false;

/**
* Whether trace managed ledger task execution time
*/
private boolean traceTaskExecution = true;

/**
* Managed ledger prometheus stats Latency Rollover Seconds
*/
private int prometheusStatsLatencyRolloverSeconds = 60;

/**
* cluster name for prometheus stats
*/
private String clusterName;
}
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,8 @@
import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedCursorInfo;
import org.apache.bookkeeper.mledger.proto.MLDataFormats.MessageRange;
import org.apache.bookkeeper.mledger.util.Futures;
import org.apache.bookkeeper.stats.NullStatsLogger;
import org.apache.bookkeeper.stats.StatsLogger;
import org.apache.bookkeeper.zookeeper.ZooKeeperClient;
import org.apache.pulsar.common.util.DateFormatter;
import org.apache.pulsar.metadata.api.MetadataStore;
Expand Down Expand Up @@ -124,15 +126,16 @@ public ManagedLedgerFactoryImpl(ClientConfiguration bkClientConfiguration, Manag

private ManagedLedgerFactoryImpl(ZooKeeper zkc, ClientConfiguration bkClientConfiguration,
ManagedLedgerFactoryConfig config) throws Exception {
this(new DefaultBkFactory(bkClientConfiguration, zkc), true /* isBookkeeperManaged */, zkc, config);
this(new DefaultBkFactory(bkClientConfiguration, zkc), true /* isBookkeeperManaged */,
zkc, config, NullStatsLogger.INSTANCE);
}

private ManagedLedgerFactoryImpl(ClientConfiguration clientConfiguration, String zkConnection, ManagedLedgerFactoryConfig config) throws Exception {
this(new DefaultBkFactory(clientConfiguration),
true,
ZooKeeperClient.newBuilder()
.connectString(zkConnection)
.sessionTimeoutMs(clientConfiguration.getZkTimeout()).build(), config);
.sessionTimeoutMs(clientConfiguration.getZkTimeout()).build(), config, NullStatsLogger.INSTANCE);
}

public ManagedLedgerFactoryImpl(BookKeeper bookKeeper, ZooKeeper zooKeeper) throws Exception {
Expand All @@ -141,22 +144,35 @@ public ManagedLedgerFactoryImpl(BookKeeper bookKeeper, ZooKeeper zooKeeper) thro

public ManagedLedgerFactoryImpl(BookKeeper bookKeeper, ZooKeeper zooKeeper, ManagedLedgerFactoryConfig config)
throws Exception {
this((policyConfig) -> bookKeeper, false /* isBookkeeperManaged */, zooKeeper, config);
this((policyConfig) -> bookKeeper, false /* isBookkeeperManaged */,
zooKeeper, config, NullStatsLogger.INSTANCE);
}

public ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory, ZooKeeper zooKeeper, ManagedLedgerFactoryConfig config)
public ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory,
ZooKeeper zooKeeper, ManagedLedgerFactoryConfig config)
throws Exception {
this(bookKeeperGroupFactory, false /* isBookkeeperManaged */, zooKeeper, config);
this(bookKeeperGroupFactory, false /* isBookkeeperManaged */, zooKeeper, config, NullStatsLogger.INSTANCE);
}

private ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory, boolean isBookkeeperManaged, ZooKeeper zooKeeper,
ManagedLedgerFactoryConfig config) throws Exception {
public ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory,
ZooKeeper zooKeeper, ManagedLedgerFactoryConfig config, StatsLogger statsLogger)
throws Exception {
this(bookKeeperGroupFactory, false /* isBookkeeperManaged */, zooKeeper, config, statsLogger);
}

private ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory,
boolean isBookkeeperManaged, ZooKeeper zooKeeper,
ManagedLedgerFactoryConfig config, StatsLogger statsLogger) throws Exception {
scheduledExecutor = OrderedScheduler.newSchedulerBuilder()
.numThreads(config.getNumManagedLedgerSchedulerThreads())
.statsLogger(statsLogger)
.traceTaskExecution(config.isTraceTaskExecution())
.name("bookkeeper-ml-scheduler")
.build();
orderedExecutor = OrderedExecutor.newBuilder()
.numThreads(config.getNumManagedLedgerWorkerThreads())
.statsLogger(statsLogger)
.traceTaskExecution(config.isTraceTaskExecution())
.name("bookkeeper-ml-workers")
.build();
cacheEvictionExecutor = Executors
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1022,6 +1022,12 @@ public class ServiceConfiguration implements PulsarConfiguration {
@FieldContext(category = CATEGORY_STORAGE_BK, doc = "Set the interval to check the need for sending an explicit LAC")
private int bookkeeperExplicitLacIntervalInMills = 0;

@FieldContext(
category = CATEGORY_STORAGE_BK,
doc = "whether expose managed ledger client stats to prometheus"
)
private boolean bookkeeperClientExposeStatsToPrometheus = false;

/**** --- Managed Ledger --- ****/
@FieldContext(
minValue = 1,
Expand Down Expand Up @@ -1224,6 +1230,19 @@ public class ServiceConfiguration implements PulsarConfiguration {
doc = "Add entry timeout when broker tries to publish message to bookkeeper.(0 to disable it)")
private long managedLedgerAddEntryTimeoutSeconds = 0;

@FieldContext(
category = CATEGORY_STORAGE_ML,
doc = "Managed ledger prometheus stats latency rollover seconds"
)
private int managedLedgerPrometheusStatsLatencyRolloverSeconds = 60;

@FieldContext(
dynamic = true,
category = CATEGORY_STORAGE_ML,
doc = "Whether trace managed ledger task execution time"
)
private boolean managedLedgerTraceTaskExecution = true;

@FieldContext(category = CATEGORY_STORAGE_ML,
doc = "New entries check delay for the cursor under the managed ledger. \n"
+ "If no new messages in the topic, the cursor will try to check again after the delay time. \n"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@

import org.apache.bookkeeper.client.BookKeeper;
import org.apache.bookkeeper.client.EnsemblePlacementPolicy;
import org.apache.bookkeeper.stats.StatsLogger;
import org.apache.pulsar.broker.ServiceConfiguration;
import org.apache.zookeeper.ZooKeeper;

Expand All @@ -35,5 +36,9 @@ BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient,
Optional<Class<? extends EnsemblePlacementPolicy>> ensemblePlacementPolicyClass,
Map<String, Object> ensemblePlacementPolicyProperties) throws IOException;

BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient,
Optional<Class<? extends EnsemblePlacementPolicy>> ensemblePlacementPolicyClass,
Map<String, Object> ensemblePlacementPolicyProperties,
StatsLogger statsLogger) throws IOException;
void close();
}
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,8 @@
import org.apache.bookkeeper.client.RackawareEnsemblePlacementPolicy;
import org.apache.bookkeeper.client.RegionAwareEnsemblePlacementPolicy;
import org.apache.bookkeeper.conf.ClientConfiguration;
import org.apache.bookkeeper.stats.NullStatsLogger;
import org.apache.bookkeeper.stats.StatsLogger;
import org.apache.commons.lang3.StringUtils;
import org.apache.pulsar.common.protocol.Commands;
import org.apache.pulsar.zookeeper.ZkBookieRackAffinityMapping;
Expand All @@ -54,7 +56,15 @@ public class BookKeeperClientFactoryImpl implements BookKeeperClientFactory {

@Override
public BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient,
Optional<Class<? extends EnsemblePlacementPolicy>> ensemblePlacementPolicyClass, Map<String, Object> properties) throws IOException {
Optional<Class<? extends EnsemblePlacementPolicy>> ensemblePlacementPolicyClass,
Map<String, Object> properties) throws IOException {
return create(conf, zkClient, ensemblePlacementPolicyClass, properties, NullStatsLogger.INSTANCE);
}

@Override
public BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient,
Optional<Class<? extends EnsemblePlacementPolicy>> ensemblePlacementPolicyClass,
Map<String, Object> properties, StatsLogger statsLogger) throws IOException {
ClientConfiguration bkConf = createBkClientConfiguration(conf);
if (properties != null) {
properties.forEach((key, value) -> bkConf.setProperty(key, value));
Expand All @@ -67,6 +77,7 @@ public BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient,
try {
return BookKeeper.forConfig(bkConf)
.allocator(PulsarByteBufAllocator.DEFAULT)
.statsLogger(statsLogger)
.build();
} catch (InterruptedException | BKException e) {
throw new IOException(e);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,12 +25,18 @@
import java.util.concurrent.RejectedExecutionException;

import org.apache.bookkeeper.client.BookKeeper;
import org.apache.bookkeeper.conf.ClientConfiguration;
import org.apache.bookkeeper.mledger.ManagedLedgerFactory;
import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig;
import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl;
import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl.BookkeeperFactoryForCustomEnsemblePlacementPolicy;
import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl.EnsemblePlacementPolicyConfig;
import org.apache.bookkeeper.stats.NullStatsProvider;
import org.apache.bookkeeper.stats.StatsLogger;
import org.apache.bookkeeper.stats.StatsProvider;
import org.apache.commons.configuration.Configuration;
import org.apache.pulsar.broker.ServiceConfiguration;
import org.apache.pulsar.broker.stats.prometheus.metrics.PrometheusMetricsProvider;
import org.apache.zookeeper.ZooKeeper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
Expand All @@ -45,6 +51,7 @@ public class ManagedLedgerClientFactory implements Closeable {
private final ManagedLedgerFactory managedLedgerFactory;
private final BookKeeper defaultBkClient;
private final Map<EnsemblePlacementPolicyConfig, BookKeeper> bkEnsemblePolicyToBkClientMap = Maps.newConcurrentMap();
private StatsProvider statsProvider = new NullStatsProvider();

public ManagedLedgerClientFactory(ServiceConfiguration conf, ZooKeeper zkClient,
BookKeeperClientFactory bookkeeperProvider) throws Exception {
Expand All @@ -56,6 +63,19 @@ public ManagedLedgerClientFactory(ServiceConfiguration conf, ZooKeeper zkClient,
managedLedgerFactoryConfig.setCacheEvictionFrequency(conf.getManagedLedgerCacheEvictionFrequency());
managedLedgerFactoryConfig.setCacheEvictionTimeThresholdMillis(conf.getManagedLedgerCacheEvictionTimeThresholdMillis());
managedLedgerFactoryConfig.setCopyEntriesInCache(conf.isManagedLedgerCacheCopyEntries());
managedLedgerFactoryConfig.setPrometheusStatsLatencyRolloverSeconds(conf.getManagedLedgerPrometheusStatsLatencyRolloverSeconds());
managedLedgerFactoryConfig.setTraceTaskExecution(conf.isManagedLedgerTraceTaskExecution());

Configuration configuration = new ClientConfiguration();
if (conf.isBookkeeperClientExposeStatsToPrometheus()) {
configuration.addProperty(PrometheusMetricsProvider.PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS,
conf.getManagedLedgerPrometheusStatsLatencyRolloverSeconds());
configuration.addProperty(PrometheusMetricsProvider.CLUSTER_NAME, conf.getClusterName());
statsProvider = new PrometheusMetricsProvider();
}

statsProvider.start(configuration);
StatsLogger statsLogger = statsProvider.getStatsLogger("pulsar_managedLedger_client");

this.defaultBkClient = bookkeeperProvider.create(conf, zkClient, Optional.empty(), null);

Expand All @@ -80,7 +100,7 @@ public ManagedLedgerClientFactory(ServiceConfiguration conf, ZooKeeper zkClient,
return bkClient != null ? bkClient : defaultBkClient;
};

this.managedLedgerFactory = new ManagedLedgerFactoryImpl(bkFactory, zkClient, managedLedgerFactoryConfig);
this.managedLedgerFactory = new ManagedLedgerFactoryImpl(bkFactory, zkClient, managedLedgerFactoryConfig, statsLogger);
}

public ManagedLedgerFactory getManagedLedgerFactory() {
Expand All @@ -91,16 +111,22 @@ public BookKeeper getBookKeeperClient() {
return defaultBkClient;
}

public StatsProvider getStatsProvider() {
return statsProvider;
}

@VisibleForTesting
public Map<EnsemblePlacementPolicyConfig, BookKeeper> getBkEnsemblePolicyToBookKeeperMap() {
return bkEnsemblePolicyToBkClientMap;
}

@Override
public void close() throws IOException {
try {
managedLedgerFactory.shutdown();
log.info("Closed managed ledger factory");

statsProvider.stop();
try {
defaultBkClient.close();
} catch (RejectedExecutionException ree) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,12 @@
import java.util.HashSet;
import java.util.Map;
import java.util.Set;
import java.io.StringWriter;
import java.io.Writer;
import java.util.Enumeration;

import org.apache.bookkeeper.stats.NullStatsProvider;
import org.apache.bookkeeper.stats.StatsProvider;
import org.apache.pulsar.broker.PulsarService;
import static org.apache.pulsar.common.stats.JvmMetrics.getJvmDirectMemoryUsed;

Expand Down Expand Up @@ -85,6 +89,8 @@ public static void generate(PulsarService pulsar, boolean includeTopicMetrics, b

generateBrokerBasicMetrics(pulsar, stream);

generateManagedLedgerBookieClientMetrics(pulsar, stream);

out.write(buf.array(), buf.arrayOffset(), buf.readableBytes());
} finally {
buf.release();
Expand Down Expand Up @@ -155,6 +161,21 @@ private static void parseMetricsToPrometheusMetrics(Collection<Metrics> metrics,
}
}

private static void generateManagedLedgerBookieClientMetrics(PulsarService pulsar, SimpleTextOutputStream stream) {
StatsProvider statsProvider = pulsar.getManagedLedgerClientFactory().getStatsProvider();
if (statsProvider instanceof NullStatsProvider) {
return;
}

try {
Writer writer = new StringWriter();
statsProvider.writeAllMetrics(writer);
stream.write(writer.toString());
} catch (IOException e) {
// nop
}
}

private static void generateSystemMetrics(SimpleTextOutputStream stream, String cluster) {
Enumeration<MetricFamilySamples> metricFamilySamples = CollectorRegistry.defaultRegistry.metricFamilySamples();
while (metricFamilySamples.hasMoreElements()) {
Expand Down
Loading