В руководстве на InfoQ Джошуа Олуикпе показал, как обрабатывать сообщения одного диалога последовательно, а разных диалогов — параллельно. Автор описывает реализацию на Go для платформы разговорного ИИ, которая использует Kafka, систему передачи и хранения сообщений.
Для каждого активного диалога создаётся свой обработчик, который обрабатывает сообщения по очереди. При временной ошибке обработчик повторяет попытку, не переходя к следующему сообщению.
Чтобы после сбоя не пропустить незавершённую работу, система фиксирует прогресс лишь до последнего подряд обработанного сообщения.
Порядок проверили в синтетическом тесте на 50 тысячах сообщений и десяти диалогах, с временными ошибками у 10% сообщений: нарушений не обнаружили. В нагрузочных тестах на высоких скоростях проверяли производительность и ошибки, но порядок сообщений отдельно не проверяли.
Проверка утверждений:
- Джошуа Олуикпе описал реализацию конвейера Kafka на Go, который сохраняет порядок сообщений внутри каждого диалога и обрабатывает разные диалоги параллельно. (подтверждено самой публикацией: доказательство; «I led the design and implementation of the core Kafka messaging layer in Go for a real-time conversational AI platform serving enterprise customers. Working closely with my team, we built a pipeline where every message flows through four processing stages: ingest, natural language understanding (NLU), large language model (LLM) inference, and delivery. The ordering requirement is strict: messages within a session are processed in sequence across all four stages, regardless of how long an individual stage takes. The system has processed over 40 million messages in production with no observed ordering violations, across thousands of concurrent chat sessions. The ordering guarantee here does not come from staying under some safe throughput. It is structural, built from consistent hashing and one goroutine per session. We have since pushed send-side throughput past 100,000 messages/second in testing with zero send errors. Architecture Overview The core challenge is that Kafka guarantees ordering within a partition, but a single partition carries interleaved traffic from many sessions. To process sessions in parallel while preserving within-session order, the application layer needs its own routing.»)
- Автор описывает реализацию на Go для платформы разговорного ИИ, которая использует Kafka. (подтверждено самой публикацией: доказательство; «I led the design and implementation of the core Kafka messaging layer in Go for a real-time conversational AI platform serving enterprise customers.»)
- Для каждого активного диалога создаётся свой обработчик, который обрабатывает сообщения по очереди. (подтверждено самой публикацией: доказательство; «Each active session is handled by its own goroutine, which processes messages one at a time in arrival order. These goroutines are created on demand and removed after a period of inactivity, keeping resource usage proportional to the number of active sessions.»)
- При временной ошибке обработчик повторяет попытку, не переходя к следующему сообщению. (подтверждено самой публикацией: доказательство; «If the failure is transient (for example, a downstream service timeout), the goroutine waits using exponential backoff with jitter before retrying the same message. Because the goroutine is blocked during the retry, later messages for that session remain queued and cannot overtake the failed message.»)
- Чтобы после сбоя не пропустить незавершённую работу, система фиксирует прогресс лишь до последнего подряд обработанного сообщения. (подтверждено самой публикацией: доказательство; «Contiguous watermark commits make concurrent processing recoverable. Offsets are committed only up to the highest contiguous completed offset, preventing a later completed message from causing an earlier unfinished message to be skipped after a crash.»)
- В синтетическом тесте на 50 тысячах сообщений и десяти диалогах, с временными ошибками у 10% сообщений, нарушений порядка не обнаружили. (подтверждено самой публикацией: доказательство; «The first validation pass used simulated downstream processing rather than the real services, with 10% of messages injected with a transient error to verify that retries did not break ordering under partial failure. The latency figures in Table 1 therefore reflect the benchmark harness and messaging path rather than real LLM inference latency. Metric Result Target Raw throughput 14,027 msg/s >10k msg/s E2E latency p50 11.7ms <200ms E2E latency p95 115ms <200ms E2E latency p99 185ms <200ms Ordering violations 0 0 Session affinity 100% 100% Table 1. Synthetic benchmark results Zero violations across 50,000 messages, 10 concurrent sessions, with 10% failure injection.»)
- В нагрузочных тестах на высоких скоростях проверяли производительность и ошибки, но порядок сообщений отдельно не проверяли. (подтверждено самой публикацией: доказательство; «The high-rate load tests measured throughput, latency, and error rate, but did not independently verify per-session ordering.»)
Первоисточники:
оценка 49,8 из 100 · тип: руководство