Table of Contents
导言:事件处理速度的迫切需要
低潜伏性应用构成了现代数字交互的支柱,每毫秒都涉及这个问题。 金融交易平台、实时欺诈检测、多人游戏游戏和IOT传感器网络都依赖于处理事件,而处理的延迟程度最小,以提供准确的响应并保持用户的信任。这些系统的核心是事件处理管道 — — 一系列阶段,它们几乎实时地摄取、过滤、转换和输出数据。优化这些管道不仅仅是一个选项;它是实现竞争优势和业务可靠性的要求。 本条探讨了事件处理管道、可操作优化策略以及持续大规模低潜伏性能所必需的持续监测纪律的核心组成部分。
了解事件处理管道
事件处理管道是运行在流数据上的处理步骤链。每个阶段都收到事件,进行特定操作,并将结果传递到下一阶段。管道的总体延迟是每个阶段所花费的时间和各阶段之间所花费的移动数据的时间之和。对于真正低延迟,每个阶段的设计必须尽量减少间接费用。
数据摄入
管道开始于摄入 — 从外部来源接收事件, 如网络服务器、 消息经纪人或硬件传感器。 摄入必须处理可变输入率和潜在的大货币。 常见技术包括Apache Kafka、 NATS、 RabbitMQ 或基于 UDP 的定制接收器。 这里的关键优化包括使用非屏蔽 I/O 、 集合连接, 以及尽可能使用零复制解析 。 例如, Kafka 的 [ [FLT: 0]] 批量压缩 [[[FLT: 2] 和 [[FLT: 2] 的memory-mapped文件[ ) 能够减少读取空闲 。
过滤
过滤会提前移除不相关的事件以减少下游处理负荷。 此阶段通常执行简单的上游检查。 为了最小化时间, 过滤会运行在事件最原始的形式( 如在完全解析前的字节上) 。 使用 [[FLT: 0]] 血压过滤器 [[[FLT: 1]] 或 [[FLT: 2]] 概率数据结构 [ 在高通量情况下可以加速成员检查 。
转变
转换可以丰富、聚合或改变事件数据。 这个阶段通常最需要计算。 常见操作包括数据格式转换、 字段提取、 窗口汇总和机器学习推论。 这里的优化包括使用 [[FLT: 0] 柱形数据模型[[[FLT: 1] 、 预分配缓冲器, 以及 [[FLT: 2]] 正在编译的表达式[。 对于汇总管道, 请考虑使用 倾覆或滑动的窗口[ , 并有高效的状态管理。
产出
最后阶段将处理过的事务传送到汇中, 如数据库、 API 或下游管道。 输出必须可靠且快。 技术包括 [ [FLT: 0]] a 同步写 [[FLT: 1] , [[FLT: 2]] 夹击 [ (同时小心的冲洗间隔以避免增加空闲性) , 以及连接集合 。 在写到数据库时, 使用已准备好的语句和索引可以减少每字的间接费用 。
优化战略
优化管道需要整体观点——一个阶段的变化影响到其他阶段,以下是具有实际执行指导的关键战略。
减少带有精益数据结构的超头处理
避免在热循环内创建对象。 重新使用可变的容器, 使用原始阵列而不是框型, 并且倾向于使用 [ [FLT: 0] 的外壳内存 [[[FLT: 1] 的数据, 用于在微型弹管上停留。 例如, 在基于 Java 的管道中, 使用 [ [[FLT: 2]] 的 FlatBuffers [[[FLT: 3] 或 [[FLT: 4] 的协议缓冲器[[[FLT: 5]] , 直接字节缓冲器可以避免堆积分配。 在 Apache Flink 等系统中, [[FLT: 6] 管理内存 [[[FLT: 7] 的功能是预分配离堆存储以减少GC 压力。
平行处理和决定货币
现代CPU架构倾向于并行主义. 将管道拆分为独立的阶段, 可以同时使用 线程集合 , 代理模型 [ (例如,Akka), 或 数据流框架 [[] (例如, Apache Flink, Kafka Streams) 。 然而,并行主义引入了命令保证和同步成本 。 使用 [ 的无锁数据结构 [ (例如, 断层环缓冲) 和 处理器, 以分解争议。 对于状态操作, 键- 以分区 确保用同一线处理事件, 不全局锁处理顺序 。
高效数据序列化
串行化往往是管道延迟的最大单一贡献者。 选择一个串行化格式, 将速度、 计划演化和互操作性相交换。 对于绝对低延迟, [[FLT: 0]] FlatBus[[[FLT: 1]] 和 [[FLT: 2] Cap 'n Proto[] 允许零复制读取数据—— 数据直接从缓冲器中获取,而不进行解码。 Apache Avro 是一个很好的选择, 当需要变换出计划时, 但需要完全解串行化。 Basure 您的串行化在现实有效载大小下, 有时会有一个简单的定制二进式超越一般用途库。 外部资源: ] Ora I/Ocle[[[9]。
优化网络通信
网络延迟通常是一个硬约束。 使用 [[ [FLT: 0]]] RDMA [[FLT: 1] 或 [[FLT: 2]] InfiniBand 进行节点间传输, 减少其延迟。 在应用程序层, 发送前进行批量事件( 但批量大小小到不增加延迟 ) 。 使用 [[FLT: 4] ] TCP NODELAY [[FLT: 5] 来禁用 Nagle 算法。 对于高频交易系统, [[[FLT: 6]] 内核绕行[[FLT: 7] 技术, 如 DPDK 或 Solarflare 的内核- bypass TCP 允许使用微秒的空间网络, 切换空格 。
调用硬件加速
GPU和FPGA在过滤和转换中常见的大规模平行计算上表现优异. 例如, Jetson GPU[可以用于实时视频分析管道,而FPGA在用于订单匹配的金融交换中很受欢迎. 然而,硬件加速增加了复杂性,最好保留在热路上. 评估CPU和加速器之间数据传输的间接费用:通常只为足够大批次实现好处.
后压和流控
非控制输入可以覆盖管道并引起空隙性悬浮. 执行回压:下游拥堵时上游阶段减速. Reactive stream(例如Project Rector Akka Streams)提供标准的后压信号. 在以卡夫卡为基础的管道中,消费者群体再平衡[和max.poll.records配置帮助控制摄入. 总是监测消费者滞后情况,作为后压问题的主要指标.
监测和培训
优化是一个不断的衡量、分析和调整周期,没有准确的监测,工作就变得盲目。
音轨密钥量表
- 端到端的延迟[(p50,p99,p999)——管道性能的最终衡量.
- 通过put——每秒事件进入和退出每个阶段.
- CPU使用和GC暂停——识别序列化瓶颈或内存压力.
- 网路往返时间和包件丢失[]——用于远程管道阶段.
- 每个阶段的队列深度——表示反压或不平衡容量.
用于剖析和可视化的工具
用于收集度量衡和用于仪表板的[PrometheusGrafana。用于分布式追踪(确定造成延迟的阶段的必要条件),]杰格或[齐普金可以通过管道追踪个别事件。async-profile为Java应用程序提供CPU的火焰图和分配热点。对于网络性能,[perf和[tcdump[13]帮助诊断内核-级别延迟。外部链接:]
教学战略
- 直接通量 :增加线程,直到CPU捆绑的操作饱和点;避免过量订阅.
- 缓冲尺寸:较大的缓冲器增加吞吐量但增加耐久性. Tune将耐久性保持在理想的p99范围内.
- 批量大小:用于写作,仅在控制了冲洗间隔时才批量;使用大小和时间的冲洗在一起.
- 枪炮收集[:在JVM管道中,切换到G1GC或ZGC,并在旧代中直接分配大对象.
- CPU pinning[:将管道线程绑定到特定核心上,可以改善缓存位置,减少上下文切换.
高级考虑
对于极低潜伏度系统,进一步的建筑模式开始发挥作用.
事件搜索和CQRS
事件源代码存储所有状态变化作为事件日志, 允许确定重放。 结合命令查询责任隔离( CQRS) , 读取模式可以优化低纬度查询, 而写入操作仍然只保留附件。 这样可以将管道与数据库瓶颈脱钩 。
国家诉无国籍人处理
无国籍阶段更容易进行规模化和优化。 但是, 许多使用的案例( 如用户会话汇总) 需要状态 。 使用 [[FLT: 0]] 嵌入状态存储 [[FLT: 1] (类似于Kafka Streams中的 RocksDB) 或 [[FLT: 2] 模拟地图 复制。 对于必须幸存失败的状态, 请考虑 RocksDB [ 或 Redis 。 使用 时间与寿命相差的迁移(TTL) 。
流程处理框架
诸如[ Apache Flink , Kafka Streams ,和[ Apache Beam[] 等框架提供了内在优化:操作员链路、状态管理、检查和精确的语义。它们抽象了许多低层次的关切问题,但增加了自己的管理费用。对于超低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低低
结论
优化低潜伏度的事件处理管道是一个多面性学科,它跨越软件设计、硬件开发以及连续的性能工程。首先要了解管道的数据流,并衡量每个阶段的当前性能。应用目标优化:精益数据结构、平行性、高效序列化以及硬件加速。永远不要停止监测;使用诸如Prometheus和Jaeger等工具及早检测回归。通过一种方法,你可以建立事件处理管道,在微秒内响应,为最要求高的应用程序解锁实时能力。关于进一步阅读,见[ Confluent在Kafka Latency和 LinkedIn的流处理架构。