流处理技术能够持续不断地处理到达的数据。它并非将数据分批收集并定期处理,而是对流经系统的每个事件进行实时分析。例如,信用卡交易的欺诈检测会在几毫秒内完成;传感器读数一旦超过阈值就会触发警报;点击流会实时更新推荐模型。整个处理过程都在动态进行,而非静止状态。
这种架构与批处理不同。数据通过 Kafka 或 Kinesis 等消息队列到达。Flink、Spark Streaming 或 Kafka Streams 等流处理器消费事件,应用转换,并发出结果。系统必须处理延迟到达的事件、乱序数据以及精确一次语义。这些都是难题。延迟到达的事件需要水印来判断时间窗口何时结束。乱序事件需要缓冲和重新排序。精确一次处理需要源、处理器和接收器之间的协调。正确实现这些非常复杂。一旦出错,结果就会不一致或重复。流处理对于需要即时洞察的用例非常强大,但对于夜间报告来说则过于复杂。批处理和流处理之间的选择取决于延迟要求。如果可以接受几分钟或几小时的延迟,批处理更简单、更经济。如果延迟时间以秒为单位,则流处理是唯一选择。
流处理特性
- 连续式——处理到达的事件
- 低延迟——毫秒至秒级
- 有状态的 — 在事件之间保持状态
- 复杂——处理延迟和乱序数据
- 可扩展——分布在多个节点上
流处理是批处理的一种反向操作。它不是等待数据累积,而是对每个出现的事件立即进行处理。
Comments
No comments yet. Be the first to share a thought.
Leave a comment