Hacktoberfest 2026:メンテナが10月に向けて印を付けた、オープンで初心者向けの issue。 Hacktoberfest の issue を見る

Advice on increasing consumer concurrency & throughput

オープン
#372 コメント 0 件 リアクション 1 件 担当者 0 名 GitHub で見る

まだ誰も着手していません。

評価

難易度
4/5
見積もり時間
3〜5日
初心者へのやさしさ
25/100
issue の種類
ドキュメント
明瞭さ
説明が足りない
活発さ
停滞
技術スタック
javascript, node.js, typescript

調査の方向性

このIssueには、Node.js Pulsar consumerのlistener APIとreceive APIを使用するTypeScriptのインライン例が2つありますが、リポジトリのファイルやテストは指定されていません。まずconsumer APIのドキュメントと既存の例を読み、次に、複数のconsumerと並行ワーカーを比較する文書化された推奨事項がプロジェクトに必要かどうかを判断してください。合意された、保守可能なベストプラクティスの説明ができれば完了です。

索引モデルが issue の本文から書いたものです。

説明

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 Shared subscription mode. Strict ordering does not matter for us
  • To keep debugging simple we have not enabled batching
  • We've been using the listener pattern

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

  1. Increases the number of consumers and uses the listener pattern
  2. 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;
    }
  };
}
主要言語
C++
スター
164
フォーク
99
平均マージ
5日 18時間
マージ済み PR(30日)
2

環境構築

  • Dockerfile・Docker Compose ファイルなし
  • プルリクエストのテンプレートあり
  • コントリビューションガイドなし

はじめの一歩

  1. issue を最後まで読み、次にプロジェクトのコントリビューションガイドを読みます。
  2. 着手することを issue にコメントします — 二人が同じ作業をするのを防げます。
  3. リポジトリをフォークし、ブランチを切って変更します。
  4. issue 番号を参照したプルリクエストを送ります。

apache/pulsar-client-node のほかの issue

apache/pulsar-client-node の issue をすべて見る

似ている issue

C++ の issue をもっと見る

新しい issue をメールで受け取る

初心者向けの GitHub issue を短くまとめたダイジェスト。