הגדרת תור הודעות: Apache Kafka למיקרוסרוויסים

ארכיטקטורת מיקרוסרוויסים דורשת העברת נתונים אמינה, ו-Apache Kafka הוא המפתח. אנו מגדירים תורי הודעות עבור הפרויקט שלכם, ומספקים הטמעה מלאה מתכנון נושאים ועד תמיכה שוטפת, תוך הבטחת שירותים יציבים וניתנים להרחבה.

פיתוח ותחזוקה של כל סוגי האתרים:

אתרי מידע או יישומי אינטרנט
אתרי תדמית, דפי נחיתה, אתרי חברה, קטלוגים מקוונים, חידונים, אתרי קידום, בלוגים, מקורות חדשות, פורטלי מידע, פורומים, אגרגטורים
אתרי מסחר אלקטרוני או יישומי אינטרנט
חנויות מקוונות, פורטלי B2B, שווקים, בורסות מקוונות, אתרי קאשבק, בורסות, פלטפורמות דרופשיפינג, מנתחי מוצרים
יישומי אינטרנט לניהול תהליכים עסקיים
מערכות CRM, מערכות ERP, פורטלים ארגוניים, מערכות ניהול ייצור, מנתחי מידע
אתרי שירות אלקטרוני או יישומי אינטרנט
פלטפורמות מודעות, בתי ספר מקוונים, בתי קולנוע מקוונים, בוני אתרים, פורטלים לשירותים אלקטרוניים, פלטפורמות אירוח וידאו, פורטלים נושאיים

אלה רק חלק מהסוגים הטכניים של אתרים שאנו עובדים איתם, ולכל אחד מהם יכולים להיות מאפיינים ופונקציונליות ספציפיים משלו, וכן ניתן להתאים אותם לצרכים ולמטרות הספציפיים של הלקוח.

השירותים שאנו מציעים
מציג 1 מתוך 1כל 2062 השירותים
הגדרת תור הודעות: Apache Kafka למיקרוסרוויסים
מורכב
~5 ימים

הכישורים שלנו:

שאלות נפוצות

העבודות האחרונות

  • פיתוח אתר חברה B2B ADVANCE
    פיתוח אתר חברה B2B ADVANCE
    1504
  • פיתוח אפליקציית ווב עבור FEEDME
    פיתוח אפליקציית ווב עבור FEEDME
    1344
  • פיתוח אתר עבור BELFINGROUP
    פיתוח אתר עבור BELFINGROUP
    1052
  • פיתוח חנות מקוונת לחברת FURNORO
    פיתוח חנות מקוונת לחברת FURNORO
    1307
  • פיתוח אפליקציית ווב עבור Enviok
    פיתוח אפליקציית ווב עבור Enviok
    1050
  • פיתוח אתר לחברת FIXPER
    פיתוח אתר לחברת FIXPER
    1033

הגדרת 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). תהליך:

  1. ניתוח עומסים ועיצוב נושאים ומחיצות של Kafka (מחיצות, גורם שכפול).
  2. פריסה באמצעות תצורת Kafka Docker compose (Docker Compose) או Kubernetes (Strimzi).
  3. תצורת יצרן עם 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: לשמירת סדר ההודעות.
  4. יישום צרכנים עם אישור ידני לאחר עיבוד.
  5. ניטור באמצעות 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 חודשים. קבל ייעוץ ליום אחד — נבחן את המשימה שלך ונציע פתרון.