Ответ
Apache Kafka — это распределённый, отказоустойчивый потоковый брокер сообщений, изначально созданный в LinkedIn. В контексте Node.js-приложений он часто используется для построения высоконагруженных, масштабируемых пайплайнов данных в реальном времени, связывая микросервисы.
Ключевые концепции:
- Топик (Topic): Именованный поток данных (например,
user-events). - Партиция (Partition): Топик делится на партиции для параллельной обработки и масштабирования.
- Производитель (Producer): Приложение (например, наш Node.js-сервис), публикующее сообщения в топик.
- Потребитель (Consumer): Приложение, подписанное на топик и читающее сообщения. Потребители объединяются в Consumer Groups для распределения нагрузки.
- Брокер (Broker): Сервер Kafka, хранящий данные.
Пример использования с библиотекой kafkajs в Node.js:
// producer.js
const { Kafka } = require('kafkajs');
const kafka = new Kafka({
clientId: 'my-node-app',
brokers: ['kafka-server1:9092', 'kafka-server2:9092']
});
const producer = kafka.producer();
async function sendEvent() {
await producer.connect();
await producer.send({
topic: 'order-created',
messages: [
{
key: 'order-123', // Ключ определяет партицию
value: JSON.stringify({ orderId: 123, amount: 99.99, userId: 'user-1' })
}
]
});
console.log('Событие отправлено в Kafka');
await producer.disconnect();
}
sendEvent().catch(console.error);
// consumer.js
const { Kafka } = require('kafkajs');
const kafka = new Kafka({
clientId: 'my-node-consumer',
brokers: ['kafka-server1:9092']
});
const consumer = kafka.consumer({ groupId: 'notification-service' });
async function runConsumer() {
await consumer.connect();
await consumer.subscribe({ topic: 'order-created', fromBeginning: false });
await consumer.run({
eachMessage: async ({ topic, partition, message }) => {
const event = JSON.parse(message.value.toString());
console.log(`Получено событие для заказа ${event.orderId}`);
// Здесь логика обработки: отправить email, уведомление и т.д.
}
});
}
runConsumer().catch(console.error);
Почему Kafka, а не RabbitMQ для Node.js? Kafka лучше подходит для сценариев с очень высокой пропускной способностью, долгосрочным хранением логов событий и потоковой обработкой через Kafka Streams или библиотеки вроде node-rdkafka.