Advice on increasing consumer concurrency & throughput
Nobody has claimed this yet.
Assessment
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Newbie friendliness
- 25/100
- Issue type
- Documentation
- Clarity
- Needs clarification
- Activity status
- Stale
- Tech stack
- javascript, node.js, typescript
- Domain
- backend, distributed-systems
Research direction
The issue provides two inline TypeScript examples using the Node.js Pulsar consumer listener and receive APIs, but names no repository files or tests. Start by reading the consumer API documentation and existing examples, then determine whether the project needs a documented recommendation comparing multiple consumers with concurrent workers; done would be an agreed, maintainable best-practice explanation.
Written by the indexing model from the issue text.
Description
Hello, we need some advice on how to increase throughput in our Pulsar consumers. Here are some details:
- We run 1 consumer per pod in a kubernetes cluster
- We run in the
Sharedsubscription mode. Strict ordering does not matter for us - To keep debugging simple we have not enabled batching
- We've been using the
listenerpattern
We've found that our messages are being processed sequentially, which leads to poor throughput. We need to speed things up a bit. What we are wondering is what is the recommended way to do so. I've attached two options we are considering below
- Increases the number of consumers and uses the listener pattern
- Uses the receiver pattern with multiple workers
We'd like to understand what the community considers best practice and why. Thank you :)
Running multiple consumers per pod
import { faker } from '@faker-js/faker';
import Pulsar from 'pulsar-client';
process.env.ENVIRONMENT = 'development';
process.env.PULSAR_SERVICE_URL = 'pulsar://localhost:6650';
const PULSAR_TOPIC = `test-${faker.string.alpha(10)}`;
const PULSAR_SUBSCRIPTION = `sub-${PULSAR_TOPIC}`;
const CONCURRENCY = 5;
const SEND_NUMBER = 10;
async function handleMessage(
message: Pulsar.Message,
consumer: Pulsar.Consumer,
): Promise<void> {
console.log('Received message: ', message.getData().toString());
await new Promise((resolve) => setTimeout(resolve, 1000));
await consumer.acknowledge(message);
}
async function main() {
const client = new Pulsar.Client({
serviceUrl: process.env.PULSAR_SERVICE_URL as string,
log: logconfig(),
messageListenerThreads: CONCURRENCY,
});
console.log('Topic: ', PULSAR_TOPIC);
console.log('Subscription: ', PULSAR_SUBSCRIPTION);
// Create the main consumer
const consumers = [];
const counter = new Map<string, number>();
const subscriptionType = 'Shared';
const ackTimeoutMs = 10_000;
const nAckRedeliverTimeoutMs = 2_000;
const batchIndexAckEnabled = false;
for (let i = 0; i < CONCURRENCY; i += 1) {
const consumer = await client.subscribe({
topic: PULSAR_TOPIC,
subscription: PULSAR_SUBSCRIPTION,
subscriptionType,
ackTimeoutMs,
nAckRedeliverTimeoutMs,
receiverQueueSize: 10,
batchIndexAckEnabled,
listener: (message, consumer) => handleMessage(message, consumer),
});
consumers.push(consumer);
}
// Send messages
const producer = await client.createProducer({ topic: PULSAR_TOPIC });
for (let i = 0; i < SEND_NUMBER; i += 1) {
const msg = `test-message-${i}`;
counter.set(msg, 0);
await producer.send({ data: Buffer.from(msg) });
}
// Sleep 20 seconds to wait for the messages to be processed
await new Promise((resolve) => setTimeout(resolve, 50000));
await producer.close();
for (const consumer of consumers) {
await consumer.close();
}
process.exit(0);
}
void main();
function logconfig() {
return (level: any, _file: any, _line: any, message: any) => {
switch (level) {
case Pulsar.LogLevel.DEBUG:
console.debug(message);
break;
case Pulsar.LogLevel.INFO:
console.info(message);
break;
case Pulsar.LogLevel.WARN:
console.warn(message);
break;
case Pulsar.LogLevel.ERROR:
console.error(message);
break;
}
};
}
Increasing concurrency per consumer
import Pulsar from 'pulsar-client';
import logger from '../../utils/logger';
process.env.ENVIRONMENT = 'development';
process.env.PULSAR_SERVICE_URL = 'pulsar://localhost:6650';
const PULSAR_TOPIC = `test-${faker.string.alpha(10)}`;
const PULSAR_SUBSCRIPTION = `sub-${PULSAR_TOPIC}`;
const CONCURRENCY = 5;
const SEND_NUMBER = 10;
async function handleMessage(
message: Pulsar.Message,
consumer: Pulsar.Consumer,
): Promise<void> {
console.log('Received message: ', message.getData().toString());
await new Promise((resolve) => setTimeout(resolve, 1000));
await consumer.acknowledge(message);
}
async function main() {
const client = new Pulsar.Client({
serviceUrl: process.env.PULSAR_SERVICE_URL as string,
log: logconfig()
});
console.log('Topic: ', PULSAR_TOPIC);
console.log('Subscription: ', PULSAR_SUBSCRIPTION);
// Create the main consumer
const consumers = [];
const counter = new Map<string, number>();
const subscriptionType = 'Shared';
const ackTimeoutMs = 10_000;
const nAckRedeliverTimeoutMs = 2_000;
const batchIndexAckEnabled = false;
const consumer = await client.subscribe({
topic: PULSAR_TOPIC,
subscription: PULSAR_SUBSCRIPTION,
subscriptionType,
ackTimeoutMs,
nAckRedeliverTimeoutMs,
receiverQueueSize: 10,
batchIndexAckEnabled,
});
await listen(
consumer,
async (consumer, message) => handleMessage(message, consumer),
CONCURRENCY,
);
// Send messages
const producer = await client.createProducer({ topic: PULSAR_TOPIC });
for (let i = 0; i < SEND_NUMBER; i += 1) {
const msg = `test-message-${i}`;
counter.set(msg, 0);
await producer.send({ data: Buffer.from(msg) });
}
// Sleep 20 seconds to wait for the messages to be processed
await new Promise((resolve) => setTimeout(resolve, 50000));
await producer.close();
await consumer.close();
process.exit(0);
}
void main();
/**
* Receive messages from a Pulsar consumer and process them concurrently.
*
* @param consumer - Pulsar consumer to receive messages from.
* @param listener - Message handler function.
* @param concurrency - Maximum number of messages to process at a time.
*/
export async function listen(
consumer: Pulsar.Consumer,
listener: (
consumer: Pulsar.Consumer,
message: Pulsar.Message,
) => Promise<void>,
concurrency = 1,
): Promise<void> {
const workers = new Array<Promise<void>>();
for (let i = 0; i < concurrency; i++) {
const worker = async () => {
for (;;) {
try {
const message = await consumer.receive();
await listener(consumer, message);
} catch (err: any) {
logger.error(`Message processing error: ${err.message}`);
}
}
};
workers.push(worker());
}
await Promise.all(workers);
}
function logconfig() {
return (level: any, _file: any, _line: any, message: any) => {
switch (level) {
case Pulsar.LogLevel.DEBUG:
console.debug(message);
break;
case Pulsar.LogLevel.INFO:
console.info(message);
break;
case Pulsar.LogLevel.WARN:
console.warn(message);
break;
case Pulsar.LogLevel.ERROR:
console.error(message);
break;
}
};
}
- Dominant language
- C++
- Stars
- 164
- Forks
- 99
- Avg merge
- 5d 18h
- Merged PRs (30d)
- 2
Contributor guide
No contributing guide indexed for this repository
First steps
- Read the whole issue, then the project's contributing guide.
- Comment on the issue to say you are picking it up — it saves two people doing the same work.
- Fork the repository and make your change on a branch.
- Open a pull request that references the issue number.
More from apache/pulsar-client-node
-
Difficulty 4/5 3-5 days Newbie friendliness 35/100
apache/pulsar-client-node#477 ·
-
公司网络无法访问海外站点,如何解决? Open
Difficulty 4/5 3-5 days Newbie friendliness 35/100
apache/pulsar-client-node#455 · 3 comments ·
-
Difficulty 4/5 3-5 days Newbie friendliness 35/100
apache/pulsar-client-node#432 · 2 reactions ·
-
Difficulty 3/5 1-2 days Newbie friendliness 38/100
apache/pulsar-client-node#431 · 3 comments ·
-
Difficulty 3/5 1-2 days Newbie friendliness 35/100
apache/pulsar-client-node#429 · 1 comment ·
All issues in apache/pulsar-client-node
Similar issues
-
Difficulty 1/5 Under an hour Newbie friendliness 90/100
AXERA-TECH/ax-llm#77 ·
-
Difficulty 1/5 Under an hour Newbie friendliness 90/100
games-on-whales/wolf#509 ·
-
Difficulty 2/5 1-3 hours Newbie friendliness 74/100
-
bug-unconfirmed
Difficulty 2/5 1-3 hours Newbie friendliness 76/100
-
Difficulty 2/5 1-3 hours Newbie friendliness 74/100
NVIDIA/cuda-samples#453 ·