הגדרת Apache Kafka עבור מיקרוסרוויסים משלבת את תור ההודעות Kafka בארכיטקטורה שלך, ומטפלת ביותר מ-100,000 הודעות בשנייה ללא אובדן נתונים. בפרויקט אחד, לקוח איבד הזמנות עקב הגדרת acks=0; לאחר מעבר ל-acks=all עם enable.idempotence=true, המסירה הפכה למובטחת וסדר האירועים נשמר. כל קבוצת צרכנים של Kafka מעבדת תת-קבוצה של מחיצות כדי להבטיח חלוקת עומסים. אנו מציגים Schema Registry לניהול גרסאות הודעות, המאפשר לשירותים שונים להתפתח בבטחה. ההגדרה שלנו, המותאמת לפרויקט שלך, נמסרת תוך 3–5 ימים עם תיעוד מלא.
למה Kafka עדיפה על RabbitMQ לעיבוד זרמים
Kafka מאחסנת הודעות לפי מדיניות שמירה (לדוגמה, 7 ימים או 100 GB), ומאפשרת למספר צרכנים להפעיל מחדש ממיקומים שונים. RabbitMQ מוחקת הודעות לאחר אישור קבלה — טובה לתורי משימות, אך לא לביקורת או להפעלה חוזרת. Kafka טובה פי 10 מ-RabbitMQ בתפוקה, וזולה פי 3 עבור נפחי נתונים גדולים.
| קריטריון | Apache Kafka | RabbitMQ |
|---|---|---|
| אחסון הודעות | לפי שמירה (זמן/גודל קבוע) | עד אישור קבלה |
| הפעלה חוזרת | נתמכת (איפוס מיקום) | לא |
| עיבוד זרמים | מובנה (Kafka Streams) | דורש כלים חיצוניים |
| תפוקה מקסימלית | מיליוני הודעות בשנייה | מאות אלפי הודעות בשנייה |
| שימוש אופייני | מקור אירועים, אנליטיקה, זרמי אירועים, לוגים | תורי משימות, RPC, התראות |
לפי תיעוד Apache Kafka, התפוקה של Kafka מגיעה ל-2–3 מיליון הודעות בשנייה על אשכול של 3 צמתים; RabbitMQ מגיעה לשיא של 300–500 אלף. ברירת המחדל לשמירה היא 7 ימים, אך אנו מתאימים ללוגיקה העסקית: לדוגמה, 30 ימים לביקורות, 100 GB ללוגים.
כיצד אנו מגדירים Kafka: מחסנית, תצורות, תהליך
אנו משתמשים ב-Confluent Kafka 7.6+ עם Schema Registry חובה. עבור PHP אנו משתמשים בספריית היצרן PHP rdkafka (librdkafka); עבור Node.js, לקוח Node.js Kafka (kafkajs). תהליך:
- ניתוח עומסים ועיצוב נושאים ומחיצות של Kafka (מחיצות, גורם שכפול).
- פריסה באמצעות תצורת Kafka Docker compose (Docker Compose) או Kubernetes (Strimzi).
- תצורת יצרן עם idempotent ו-acks=all, בתוספת
# docker-compose.yml services: zookeeper: image: confluentinc/cp-zookeeper:7.6.0 environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 volumes: - zookeeper_data:/var/lib/zookeeper/data - zookeeper_log:/var/lib/zookeeper/log kafka: image: confluentinc/cp-kafka:7.6.0 depends_on: [zookeeper] environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false" KAFKA_LOG_RETENTION_HOURS: 168 # 7 дней KAFKA_LOG_RETENTION_BYTES: 107374182400 # 100 GB KAFKA_NUM_PARTITIONS: 6 KAFKA_DEFAULT_REPLICATION_FACTOR: 1 volumes: - kafka_data:/var/lib/kafka/data ports: - "9092:9092" kafka-ui: image: provectuslabs/kafka-ui:latest depends_on: [kafka] environment: KAFKA_CLUSTERS_0_NAME: local KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:9092 ports: - "8080:8080" volumes: zookeeper_data: zookeeper_log: kafka_data:לשמירת סדר ההודעות. - יישום צרכנים עם אישור ידני לאחר עיבוד.
- ניטור באמצעות JMX + Grafana עם התראות בפיגור > 10,000.
הגדרת Docker
---
# docker-compose.yml
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.6.0
environment:
ZOOKEEPER_CLIENT_PORT: 2181
ZOOKEEPER_TICK_TIME: 2000
volumes:
- zookeeper_data:/var/lib/zookeeper/data
- zookeeper_log:/var/lib/zookeeper/log
kafka:
image: confluentinc/cp-kafka:7.6.0
depends_on: [zookeeper]
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false"
KAFKA_LOG_RETENTION_HOURS: 168 # 7 дней
KAFKA_LOG_RETENTION_BYTES: 107374182400 # 100 GB
KAFKA_NUM_PARTITIONS: 6
KAFKA_DEFAULT_REPLICATION_FACTOR: 1
volumes:
- kafka_data:/var/lib/kafka/data
ports:
- "9092:9092"
kafka-ui:
image: provectuslabs/kafka-ui:latest
depends_on: [kafka]
environment:
KAFKA_CLUSTERS_0_NAME: local
KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:9092
ports:
- "8080:8080"
volumes:
zookeeper_data:
zookeeper_log:
kafka_data:
יצירת נושא
# Создать топик с 6 партициями и репликой 1 (для одиночного брокера)
kafka-topics.sh --bootstrap-server kafka:9092 \
--create \
--topic user-events \
--partitions 6 \
--replication-factor 1 \
--config retention.ms=604800000 \
--config cleanup.policy=delete
# Просмотр
kafka-topics.sh --bootstrap-server kafka:9092 --describe --topic user-events
PHP: יצרן (librdkafka)
use RdKafka\Producer;
use RdKafka\Conf;
class KafkaProducer
{
private Producer $producer;
public function __construct()
{
$conf = new Conf();
$conf->set('bootstrap.servers', config('kafka.brokers'));
$conf->set('security.protocol', 'PLAINTEXT');
$conf->set('acks', 'all'); // подтверждение от всех реплик
$conf->set('retries', '3');
$conf->set('enable.idempotence', 'true'); // ровно одна запись
$conf->set('compression.type', 'snappy');
$conf->setDrMsgCb(function ($kafka, $message) {
if ($message->err !== RD_KAFKA_RESP_ERR_NO_ERROR) {
Log::error('Kafka delivery failed', [
'error' => $message->errstr(),
'topic' => $message->topic_name,
]);
}
});
$this->producer = new Producer($conf);
}
public function publish(string $topic, string $key, array $payload): void
{
$rdTopic = $this->producer->newTopic($topic);
$rdTopic->produce(
partition: RD_KAFKA_PARTITION_UA, // автовыбор партиции по key
msgflags: 0,
payload: json_encode($payload),
key: $key, // один ключ → одна партиция → порядок событий
);
$this->producer->poll(0);
}
public function flush(): void
{
$result = $this->producer->flush(10000); // 10 секунд таймаут
if (RD_KAFKA_RESP_ERR_NO_ERROR !== $result) {
throw new \RuntimeException('Kafka flush failed: ' . rd_kafka_err2str($result));
}
}
}
// Использование
$producer->publish('user-events', (string) $user->id, [
'event' => 'user.registered',
'user_id' => $user->id,
'email' => $user->email,
'timestamp' => now()->toIso8601String(),
]);
$producer->flush();
Node.js: יצרן וצרכן על kafkajs
import { Kafka, CompressionTypes } from 'kafkajs';
const kafka = new Kafka({
clientId: 'myapp-api',
brokers: [process.env.KAFKA_BROKERS!],
retry: {
retries: 5,
initialRetryTime: 300,
factor: 0.2,
},
});
// Producer
const producer = kafka.producer({
allowAutoTopicCreation: false,
idempotent: true,
maxInFlightRequests: 5,
});
await producer.connect();
await producer.send({
topic: 'user-events',
compression: CompressionTypes.Snappy,
messages: [
{
key: String(userId),
value: JSON.stringify({
event: 'user.login',
userId,
ip,
timestamp: Date.now(),
}),
headers: {
'content-type': 'application/json',
},
},
],
});
// Consumer
const consumer = kafka.consumer({
groupId: 'audit-service',
});
await consumer.connect();
await consumer.subscribe({
topic: 'user-events',
fromBeginning: false,
});
await consumer.run({
eachMessage: async ({ topic, partition, message }) => {
const payload = JSON.parse(message.value!.toString());
await AuditLog.create({
event: payload.event,
userId: payload.userId,
metadata: payload,
});
},
});
השגת יצרן Idempotent
הפעל # Создать топик с 6 партициями и репликой 1 (для одиночного брокера) kafka-topics.sh --bootstrap-server kafka:9092 \ --create \ --topic user-events \ --partitions 6 \ --replication-factor 1 \ --config retention.ms=604800000 \ --config cleanup.policy=delete # Просмотр kafka-topics.sh --bootstrap-server kafka:9092 --describe --topic user-events ו-use RdKafka\Producer; use RdKafka\Conf; class KafkaProducer { private Producer $producer; public function __construct() { $conf = new Conf(); $conf->set('bootstrap.servers', config('kafka.brokers')); $conf->set('security.protocol', 'PLAINTEXT'); $conf->set('acks', 'all'); // подтверждение от всех реплик $conf->set('retries', '3'); $conf->set('enable.idempotence', 'true'); // ровно одна запись $conf->set('compression.type', 'snappy'); $conf->setDrMsgCb(function ($kafka, $message) { if ($message->err !== RD_KAFKA_RESP_ERR_NO_ERROR) { Log::error('Kafka delivery failed', [ 'error' => $message->errstr(), 'topic' => $message->topic_name, ]); } }); $this->producer = new Producer($conf); } public function publish(string $topic, string $key, array $payload): void { $rdTopic = $this->producer->newTopic($topic); $rdTopic->produce( partition: RD_KAFKA_PARTITION_UA, // автовыбор партиции по key msgflags: 0, payload: json_encode($payload), key: $key, // один ключ → одна партиция → порядок событий ); $this->producer->poll(0); } public function flush(): void { $result = $this->producer->flush(10000); // 10 секунд таймаут if (RD_KAFKA_RESP_ERR_NO_ERROR !== $result) { throw new \RuntimeException('Kafka flush failed: ' . rd_kafka_err2str($result)); } } } // Использование $producer->publish('user-events', (string) $user->id, [ 'event' => 'user.registered', 'user_id' => $user->id, 'email' => $user->email, 'timestamp' => now()->toIso8601String(), ]); $producer->flush(); . אנו גם מגדירים import { Kafka, CompressionTypes } from 'kafkajs'; const kafka = new Kafka({ clientId: 'myapp-api', brokers: [process.env.KAFKA_BROKERS!], retry: { retries: 5, initialRetryTime: 300, factor: 0.2, }, }); // Producer const producer = kafka.producer({ allowAutoTopicCreation: false, idempotent: true, maxInFlightRequests: 5, }); await producer.connect(); await producer.send({ topic: 'user-events', compression: CompressionTypes.Snappy, messages: [{ key: String(userId), value: JSON.stringify({ event: 'user.login', userId, ip, timestamp: Date.now() }), headers: { 'content-type': 'application/json' }, }], }); // Consumer const consumer = kafka.consumer({ groupId: 'audit-service' }); await consumer.connect(); await consumer.subscribe({ topic: 'user-events', fromBeginning: false }); await consumer.run({ eachMessage: async ({ topic, partition, message }) => { const payload = JSON.parse(message.value!.toString()); await AuditLog.create({ event: payload.event, userId: payload.userId, metadata: payload, }); }, }); כדי להבטיח סדר הודעות. היצרן מקבל מזהה ייחודי, והברוקר מבטל כפילויות. זה חובה עבור עסקאות פיננסיות וביקורות.
ניטור פיגור קבוצת צרכנים
נטר פיגור קבוצת צרכנים באמצעות enable.idempotence=true. לאוטומציה, פרוס JMX Exporter המזין את Prometheus, ולאחר מכן בנה לוחות מחוונים של Grafana עם סף התראה בפיגור > 10,000.
מה כלול בהגדרה סוהר (תוצרים מסחריים)
- ביקורת ארכיטקטורה נוכחית – ניתוח עומסים, עיצוב נושאים ומחיצות.
- פריסת אשכול – Docker Compose או Kubernetes (Strimzi) עם סובלנות לתקלות.
- שילוב יישומים – יצרנים/צרכנים ב-PHP או Node.js עם טיפול בשגיאות ולוגיקת ניסיון חוזר.
- Schema Registry – יישום סכמת Avro לניהול גרסאות.
- ניטור – לוח מחוונים של Grafana עם התראות על פיגור ועומס ברוקר.
- הדרכת צוות – תיעוד, runbooks, ונהלי שחזור מאסון.
- תמיכה לאחר השקה – 14 ימים של תגובה לאירועים וכיוונון עדין.
המדדים והתמחור שלנו
אמינים על ידי 50+ לקוחות מרוצים עם ניסיון של 5+ שנים ו-20+ פרויקטי תורים מוצלחים. זה מבטיח אמינות גבוהה ויעילות עלות. לדוגמה, לקוח אחד חוסך $15,000 בשנה בעלויות תשתית לאחר מעבר משירות מנוהל. חיסכון שנתי אופייני: $5,000–$20,000 בהשוואה לשירותי ענן מנוהלים על פני 12 חודשים. עבור פרויקט טיפוסי, העלות הכוללת נעה בין $2,500 ל-$4,000 עם חיסכון שנתי ממוצע של $10,000.
| שלב | משך | עלות משוערת |
|---|---|---|
| אשכול בסיסי + יצרן/צרכן (שפה אחת) | 3–4 ימים | החל מ-$2,500 |
| Schema Registry + Avro | +2 ימים | החל מ-$1,500 |
| Kafka Streams לצבירה | 3–5 ימים | החל מ-$3,000 |
| אשכול של 3 צמתים ב-Kubernetes | 4–5 ימים | החל מ-$4,000 |
העלות מחושבת באופן אישי; מקרה טיפוסי חוסך ללקוחות 30–50% בהשוואה לשירותים מנוהלים על פני 12 חודשים. קבל ייעוץ ליום אחד — נבחן את המשימה שלך ונציע פתרון.







