バッチ処理はストリーミングの特殊なケースです。ビジネスが時間単位ではなく秒単位で反応する必要がある場合、継続的なデータフローのために構築されたアーキテクチャが必要です。

ダッシュボードは、誰かがそれを見る頃には既に古くなっています。不正検出は夜間のバッチジョブとして実行され、翌朝に不正を検知します。在庫数は1時間ごとに更新され、過剰販売を引き起こします。センサーデータは収集されますが、夜間の ETL で分析されるまで利用されません。データがソースから処理を経てコンシューマーまで、サブ秒のレイテンシーで継続的に流れるシステムが必要です — リアルタイムアナリティクス、ライブ通知、ストリーミング AI 推論、システム間の即時同期などです。
リアルタイムストリーミングアーキテクチャは、データを離散的なバッチではなく、継続的で無制限のフローとして処理します。イベントプロデューサーはストリーミングプラットフォーム(Kafka、Kinesis、Pulsar)に発行します。ストリームプロセッサー(Flink、Kafka Streams、カスタムコンシューマー)は、イベントをインフライトで変換、エンリッチ、フィルタリング、集計します。処理された結果は、リアルタイムダッシュボード(WebSocket)、検索インデックス(Elasticsearch)、アナリティクスデータベース(ClickHouse)、およびダウンストリームサービスなどのコンシューマーにプッシュされます。CDC(Change Data Capture)により、既存のデータベースはアプリケーションの変更なしにイベントソースとして参加できます。
Explore more design patterns and system architectures
MicrocosmWorksは、マルチコンシューマリプレイ、長期間のデータ保持、およびクロスクラウドポータビリティを必要とするチームにKafkaを推奨します。これは、そのログベースのアーキテクチャが無制限のコンシューマグループが同じデータストリームを個別に再読み込みすることをサポートしているためです。Kinesisは、AWSエコシステムと密接に統合されたフルマネージドサービスを希望し、データ保持期間が7日未満で、コンシューマアプリケーションが10未満である場合に、より良い選択肢です。当社は、適切な推奨を行うため、アーキテクチャ評価中にお客様の特定の要件(スループット、データ保持期間、コンシューマパターン、運用成熟度)を評価します。
MicrocosmWorksは、冪等なプロデューサー、トランザクションコンシューマー、およびRedisのような高速ルックアップキャッシュに保存されたイベントフィンガープリントを使用する重複排除レイヤーの組み合わせを通じて、exactly-onceセマンティクスを実装しています。Kafkaベースのシステムの場合、私たちはコンシューマーオフセットとプロデューサーの書き込みをアトミックにコミットするKafkaの組み込みトランザクションAPIを活用しています。一方、カスタムストリーミングパイプラインの場合、コンシューマー側での重複排除を伴うアウトボックスパターンを実装しています。私たちは常にコンシューマーが冪等になるようにセーフティネットとして設計しています。そのため、たとえexactly-onceメカニズムがエッジケースで失敗しても、イベントを再処理しても同じ結果が生成されます。
MicrocosmWorksは通常、インジェスト、処理、シンク書き込みを含むストリーミングパイプラインに対して、50〜200msのエンドツーエンドのレイテンシを実現します。Apache FlinkやKafka Streamsのようなインメモリのストリームプロセッサを使用する、よりシンプルなパススルーやフィルタリングのワークロードでは、10ms未満も達成可能です。最大のレイテンシ要因は通常、ネットワークホップ、シリアライゼーションオーバーヘッド、およびシンク書き込みのバッチ処理であり、これらはお客様のレイテンシとスループットのトレードオフに関するご要望に基づいて調整します。アーキテクチャ設計の際、パイプラインのステージごとに明示的なレイテンシSLOを設定し、本番環境におけるp50、p95、p99のレイテンシを追跡するモニタリングダッシュボードを構築します。
MicrocosmWorksは、既存のコンシューマーを破壊することなくプロデューサーがデータ形式を進化させられるよう、後方互換性および前方互換性のルールを強制するスキーマレジストリ(通常、Confluent Schema Registry または AWS Glue Schema Registry)を実装しています。当社では、明示的なスキーマバージョニングを備えたAvroまたはProtobufシリアライゼーションを使用しており、これにより各メッセージは自己記述的となり、生成されてからスキーマが変更された場合でもデシリアライズ可能です。当社のCI/CDパイプラインには、提案されたスキーマ変更がダウンストリームコンシューマーを破壊する可能性があれば、デプロイメントをブロックする自動化されたスキーマ互換性チェックが含まれています。
MicrocosmWorks は、本番ストリーミングプラットフォームを確実に維持するために、分散システム、ストリーム処理フレームワーク、インフラ自動化の経験を持つ最低2〜3名のエンジニアを推奨しています。この専門知識を社内で構築したくない企業向けには、当社のチームがクラスター運用、パフォーマンスチューニング、インシデント対応を処理し、お客様の開発者がストリーム処理アプリケーションの構築に集中できるように、1時間あたり$15〜$40でマネージドストリーミングプラットフォームサポートを提供しています。また、4〜8週間の期間にわたり、既存のエンジニアリングチームがKafka、Flink、またはKinesisの運用に関するスキルアップを図るトレーニングプログラムも提供しています。
このアーキテクチャは4つのレイヤーで構成されています。イベントソースは、アプリケーションイベント、データベース CDC ストリーム、IoT テレメトリー、ユーザーのクリックストリーム、外部 API Webhook などのデータを生成します。ストリーミングプラットフォーム(Kafka)は、耐久性があり、順序付けられ、リプレイ可能なイベントストレージを提供します。ストリームプロセッサーは、トピックからコンシュームし、変換(フィルタリング、エンリッチメント、ウィンドウ集計、結合)を適用し、出力トピックまたはシンクに生成します。コンシューマーは、処理されたストリームをサブスクライブします — WebSocket サーバーはブラウザーにプッシュし、コネクターはデータベースにシンクし、アラートエンジンはルールを評価して通知を発行します。
| レイヤー | テクノロジー |
|---|---|
| ストリーミング | Apache Kafka (MSK, Confluent), Kinesis, Apache Pulsar, Redpanda |
| CDC | Debezium, AWS DMS, Maxwell |
| 処理 | Apache Flink, Kafka Streams, Benthos, カスタムコンシューマー |
| リアルタイム配信 | WebSocket (Socket.io), SSE, GraphQL Subscriptions |
| アナリティクス | ClickHouse, Apache Druid, Elasticsearch, TimescaleDB |
| 可観測性 | Kafka ラグ監視(Burrow)、Flink メトリクス、カスタムレイテンシートラッキング |
| 利用すべき時 | 避けるべき時 |
|---|---|
| ビジネス上の意思決定にサブ秒のデータ鮮度が必要な場合(不正、監視、取引など) | 時間単位/日単位の鮮度を持つバッチ処理でビジネスニーズが満たされる場合 |
| 複数のコンシューマーが同じイベントストリームを必要とする場合(ファンアウト、疎結合システム) | 単一のプロデューサーと単一のコンシューマーの場合 — シンプルなキューで十分 |
| デバッグ、再処理、または新しいコンシューマー構築のためにイベントリプレイが必要な場合 | データ量が少なく(1Kイベント/分未満)、ストリーミングインフラストラクチャを正当化できない場合 |
| コード変更なしで既存のデータベースをダウンストリームシステムに同期するために CDC が必要な場合 | チームが分散システムの経験を欠いている場合 — ストリーミングは運用上の複雑さを大幅に増加させます |
MW は、「リプレイ原則」に基づきストリーミングシステムを設計しています — すべてのストリームは特定の時点からリプレイ可能であるべきであり、新しいコンシューマーが過去のデータをバックフィルしたり、既存のコンシューマーがバグ修正後に再処理したりできるようにします。当社の Kafka デプロイメントには、スキーマ進化ポリシー(デフォルトで後方互換)、コンシューマーラグアラート(ビジネスに影響する遅延になる前)、および自動再試行機能を備えたデッドレタートピックが含まれています。当社は、ビデオアナリティクス、IoT テレメトリー、リアルタイムダッシュボード向けに、50万イベント/秒以上を処理するストリーミングパイプラインを構築してきました。
単一のコードベース、数百のテナント、データ漏洩ゼロ — すべてのスケーラブルなSaaSビジネスの基盤。