ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Ray 事件导出基础设施深入解析:从 C++ 事件采集到 HTTP 外部服务发布

Ray 事件导出基础设施深入解析:从 C++ 事件采集到 HTTP 外部服务发布 Ray 事件导出基础设施深入解析从 C 事件采集到 HTTP 外部服务发布【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/rayRayAI compute engine内部运行着大量分布式组件如何把这些组件GCS、Worker 等产生的事件任务生命周期、Actor 状态、节点变更等统一采集、聚合并导出到外部监控或审计系统是集群可观测性的关键一环。本文基于当前仓库 Ray 2.52.1 版本系统讲解 Ray Event Exporter 基础设施事件从 C 侧创建、缓冲、合并经 gRPC 上报到 Python 侧 AggregatorAgent 聚合缓冲最终过滤、转 JSON 并通过 HTTP POST 发布到外部服务的完整链路同时给出新增自定义事件类型的逐步实践方法。架构总览事件的多级流水线Ray 的事件系统是一条多级流水线贯穿 C 组件、Python Agent 与外部服务C 组件GCS、Worker 进程通过实现RayEventInterface创建事件。当前实现中 raylet 并不产生 Ray 事件但从技术上讲并不存在阻止它产生事件的限制。事件缓冲事件在 C 侧被放入一个有界环形缓冲bounded circular buffer中暂存。事件合并在导出前相同 entity ID 与相同类型的事件会被合并以减少数据量。gRPC 导出事件通过 gRPC 从 C 组件导出到聚合 AgentAggregatorAgent。Python 聚合AggregatorAgent接收事件并将其放入MultiConsumerEventBuffer缓冲。HTTP 发布事件经过过滤、转换为 JSON 后发布到外部 HTTP 服务。整体数据流如下C Components (GCS, workers) ↓ (通过 RayEventInterface 创建事件) RayEventRecorder (C) ↓ (缓冲 合并事件) ↓ (通过 EventAggregatorClient 进行 gRPC 导出) AggregatorAgent (Python) ↓ (加入 MultiConsumerEventBuffer) RayEventPublisher ↓ (过滤 转换为 JSON) ↓ (HTTP POST) 外部 HTTP 服务其中 C 侧的入口实现在 ray_event_recorder.h、事件抽象在 ray_event_interface.hPython 侧的聚合与发布分别在 aggregator_agent.py 与 ray_event_publisher.py。事件类型与结构Ray 事件使用 protobuf 消息描述结构以基础消息RayEvent作为外壳内部嵌套各事件类型专属的消息体。事件类型EventType事件按类型划分定义在 events_base_event.proto 的EventType枚举中。以 2.52.1 版本仓库为准当前包含如下类型事件类型说明TASK_DEFINITION_EVENT任务定义信息TASK_LIFECYCLE_EVENT任务状态迁移同时覆盖普通任务与 Actor 任务ACTOR_TASK_DEFINITION_EVENTActor 任务定义TASK_PROFILE_EVENT任务 profiling 数据DRIVER_JOB_DEFINITION_EVENTDriver job 定义DRIVER_JOB_LIFECYCLE_EVENTDriver job 状态迁移NODE_DEFINITION_EVENT节点定义NODE_LIFECYCLE_EVENT节点状态迁移ACTOR_DEFINITION_EVENTActor 定义ACTOR_LIFECYCLE_EVENTActor 状态迁移PLACEMENT_GROUP_DEFINITION_EVENT放置组Placement Group定义PLACEMENT_GROUP_LIFECYCLE_EVENT放置组状态迁移SUBMISSION_JOB_DEFINITION_EVENT提交型 JobSubmission Job定义SUBMISSION_JOB_LIFECYCLE_EVENT提交型 Job 状态迁移PLATFORM_EVENT平台事件WORKER_LIFECYCLE_EVENTWorker 状态迁移WORKER_DEFINITION_EVENTWorker 定义说明文档原列表以 10 种核心类型为例从当前仓库的EventType枚举events_base_event.proto#L56-L75可见事件类型已扩充到 17 种新增了 Placement Group、Submission Job、Platform、Worker 等类型。新增类型时需要注意枚举值不能与现有值冲突。事件结构RayEvent 消息基础RayEvent消息events_base_event.proto#L39-L130包含以下公共字段event_idbytes事件的唯一标识符source_typeSourceType产生事件的组件。枚举定义包括CORE_WORKER、GCS、RAYLET、CLUSTER_LIFECYCLE、AUTOSCALER、JOBSevent_typeEventType事件类型。该字段的意义在于无需反序列化嵌套消息即可判断事件类型timestampgoogle.protobuf.Timestamp事件创建时的 epoch 时间戳severitySeverity事件严重级别。枚举为TRACE、DEBUG、INFO默认、WARNING、ERROR、FATALmessagestring可选的字符串消息session_namestring当前 Ray session 标识node_idbytes产生事件的节点 ID嵌套事件消息针对各事件类型的专属消息例如task_definition_event、actor_lifecycle_event等。每个RayEvent消息中只期望设置其中一个字段。实体 IDEntity ID概念实体 IDentity ID是与事件关联实体的唯一标识它的作用有两个关联Association将执行事件与定义事件关联起来。例如任务生命周期事件与任务定义事件通过同一实体 ID 建立联系合并Merging在导出前将相同实体 ID 与相同类型的事件分组合并从而压缩数据量。实体 ID 的取值规则在 ray_event_interface.h#L29-L38 中有明确注释例如任务的实体 ID 是task_id task_attempt任务 ID 与任务尝试 ID 的组合Driver job 的实体 ID 是 job ID。其他典型映射如下任务事件task_id task_attempt作为实体 IDActor 事件actor_id作为实体 IDDriver job 事件job_id作为实体 ID。C 侧事件记录与缓冲RayEventRecorderC 组件通过RayEventRecorder类记录事件它提供线程安全的事件缓冲与导出能力。类定义位于 ray_event_recorder.h。RayEventRecorder 的职责RayEventRecorder继承自RayEventRecorderBase是一个线程安全的记录器具备以下能力维护一个有界环形缓冲boost::circular_buffer见 ray_event_recorder.h#L52存储待导出事件在导出前将相同实体 ID 与相同类型的事件合并通过EventAggregatorClient周期性地将事件经 gRPC 导出到聚合 Agent缓冲满时跟踪被丢弃的事件数量。从类注释可见RayEventRecorderBase统一负责导出循环与 gRPC 排空drain机制各子类只需实现AddEvents与ExportEventsray_event_recorder_base.h#L33-L37。添加事件AddEvents事件通过AddEvents()方法加入记录器ray_event_recorder.cc#L65-L84该方法接收一个RayEventInterface指针的 vector。处理流程为检查是否启用先检查记录器自身的enabled_标志再检查RayConfig::instance().enable_ray_event()配置对应 Ray 配置项enable_ray_event默认false见 ray_config_def.h#L701。未启用时直接返回不产生任何开销计算容量判断data_list.size() buffer_.size()是否会超过max_buffer_size_丢弃旧事件并记录指标若超限计算需要移除的事件数打印Dropping N events from the buffer.错误日志并通过dropped_events_counter指标记录丢弃数量携带Source标签标明来源组件名入缓冲将新事件逐个push_back进环形缓冲。这里使用的是absl::Mutex加锁确保多线程并发添加时数据安全AddEvents与ExportEvents共享同一把mutex_。缓冲管理记录器使用boost::circular_buffer存储事件有界容量缓冲满时最旧的事件会被丢弃以腾出空间给新事件丢弃跟踪被丢弃的事件通过dropped_events_counter指标统计指标携带来源组件名标签默认大小与配置默认缓冲大小为 10,000 个事件可通过RAY_ray_event_recorder_max_queued_events环境变量配置。对应 Ray 配置项定义在 ray_config_def.h#L703在 gcs_server.cc#L174 等位置被实际使用。C 侧事件导出gRPCC 组件通过 gRPC 将事件导出到聚合 Agent导出流程由调用StartExportingEvents()启动。StartExportingEventsStartExportingEvents()位于RayEventRecorderBaseray_event_recorder_base.h#L42其职责包括检查事件记录是否启用校验是否已被调用过该方法只能调用一次后续调用会被忽略设置一个PeriodicalRunner周期性调用ExportEvents()使用配置的导出间隔ray_events_report_interval_ms默认1000ms见 ray_config_def.h#L616。相应地基类还提供StopExportingEvents()在优雅关闭时执行一次最终 flush确保缓冲中的事件在退出前都被发送出去ray_event_recorder_base.h#L44-L46。ExportEvents 处理流程ExportEvents()的具体实现在 ray_event_recorder.cc#L39-L63检查缓冲加锁后若缓冲为空直接返回防重入检查grpc_in_progress_标志若上一次 gRPC 导出仍在进行中打印告警并返回避免重叠请求。这一步同时防止StopExportingEvents()与进行中的周期导出产生竞态排空缓冲将缓冲中的事件逐个移入一个临时std::list然后清空缓冲分组与序列化调用共享的GroupAndSerializeEvents()按entity_id, event_type分组对每组调用合并再序列化为RayEventprotobuf发送通过EventAggregatorClient::AddEvents()经 gRPC 发送给聚合 Agent清除缓冲成功导出后缓冲已清空。事件合并逻辑Merging事件合并在GroupAndSerializeEvents()中完成按实体 ID 与类型分组然后调用各事件实现类的Merge()方法接口定义见 ray_event_interface.h#L40-L56。合并是一种针对数据体积的优化定义类事件Definition Events合并时通常不发生变化例如 Actor 定义生命周期类事件Lifecycle Events状态迁移会被追加形成一条时间序列。例如任务状态迁移started → running → completed会被合并进同一个事件中。接口注释中给出了直观示例三个{entity_id: 1, type: task, state_transitions: [...]}事件分别携带(started, 1000)、(running, 1001)、(completed, 1002)合并后成为一个携带完整状态迁移序列的事件。另外接口还定义了SupportsMerge()ray_event_interface.h#L70不支持合并的事件类型当前从 Python 侧发送的事件不支持合并返回false记录器会单独发送它们而不按entity_id, event_type分组。错误处理若 gRPC 导出失败打印错误日志进程继续运行不会崩溃下一个导出周期会再次尝试发送事件会保留在缓冲中直到成功导出或缓冲已满、旧事件被丢弃。Python 侧事件接收与缓冲AggregatorAgentAggregatorAgent是 dashboard agent 模块aggregator_agent.py#L85-L93负责通过 gRPC 服务接收 C 组件发来的事件并进行缓冲供发布器使用。其核心职责包括实现EventAggregatorServiceServicer提供 gRPC 事件接收能力维护MultiConsumerEventBuffer作为事件存储管理RayEventPublisher实例向外部 HTTP 端点发布事件跟踪事件接收、缓冲与发布相关指标。AddEvents gRPC 处理器AddEvents()aggregator_agent.py#L214-L253是接收事件的 gRPC 处理器检查事件处理是否启用_event_processing_enabled未启用时直接返回空回复记录请求中的事件总数received_count若启用了向 GCS 或 dashboard head 的发布先把task_events_metadata合并进对应的 metadata buffer遍历请求中的每个事件逐个调用MultiConsumerEventBuffer.add_event()写入缓冲处理失败单个事件添加失败时计数failed_count并记录错误日志日志中含事件 ID记录指标prefix_events_received_total接收总数与prefix_events_buffer_add_failures_total添加失败数。MultiConsumerEventBufferMultiConsumerEventBuffermulti_consumer_event_buffer.py#L26-L36是一个 asyncio 友好的事件缓冲特性如下多消费者支持每个消费者consumer拥有独立的游标索引cursor index。RayEventPublisher与其他消费者共享同一个缓冲驱逐跟踪Evictions缓冲满时最旧的事件被丢弃并按消费者分别跟踪被驱逐的事件数有界缓冲内部使用deque(maxlenmax_size)限制缓冲大小asyncio 安全使用asyncio.Lock与asyncio.Condition做同步。关键操作add_event()multi_consumer_event_buffer.py#L62-L94向缓冲添加事件满时丢弃最旧事件。若被丢弃的事件恰好是某消费者下一个要消费的事件则记录该消费者的驱逐指标否则将该消费者的游标减一进行修正wait_for_batch()multi_consumer_event_buffer.py#L114-L177等待一批事件最多max_batch_size个。等待分两个阶段第一阶段无限期等待直到至少有一个事件可消费保证返回的批次至少含一个事件第二阶段在超时时间内尽量攒满批次。因此该方法的超时参数只在缓冲中已存在事件时生效register_consumer()multi_consumer_event_buffer.py#L179-L190以唯一名称注册新消费者重复注册会抛出ValueError。缓冲大小与批量大小在 aggregator_agent.py#L44-L57 中定义RAY_DASHBOARD_AGGREGATOR_AGENT_MAX_EVENT_BUFFER_SIZE默认 100,000该默认值从 1,000,000 下调而来因为此前观察到 1,000,000 条事件可能占用高达 20 GB 内存RAY_DASHBOARD_AGGREGATOR_AGENT_MAX_EVENT_SEND_BATCH_SIZE默认 1,000。事件过滤Exposable Event Types聚合 Agent 通过_can_expose_event()判断事件是否可以暴露给外部服务。该逻辑实现在 async_publisher_client.py#L72-L84只有事件类型在可暴露类型集合exposable event types中才允许被发布。可暴露类型默认集合定义在 configs.py#L34-L41注意TASK_PROFILE_EVENT默认不暴露给外部服务。事件发布到 HTTPRayEventPublisher事件由RayEventPublisher发布到外部 HTTP 服务它从事件缓冲读取批次并发送 HTTP POST 请求。类实现位于 ray_event_publisher.py。RayEventPublisher 工作循环RayEventPublisherray_event_publisher.py#L54-L58运行一个run_forever()工作循环注册为MultiConsumerEventBuffer的消费者通过wait_for_batch()持续等待事件批次等待窗口由PUBLISHER_MAX_BUFFER_SEND_INTERVAL_SECONDS控制默认 0.1 秒见 configs.py#L24-L26使用配置的PublisherClientInterface发布批次失败时按指数退避重试记录发布成功、失败、延迟等指标。发布器运行在异步上下文中全程使用asyncio进行非阻塞操作。若目标未启用例如未配置 HTTP 端点则使用NoopPublisherray_event_publisher.py#L274-L288运行但不做任何事。AsyncHttpPublisherClientAsyncHttpPublisherClientasync_publisher_client.py#L97-L122负责实际的 HTTP 发布事件过滤使用events_filter_fn典型实现即_can_expose_event过滤事件。若配置为ALL则放行所有事件类型JSON 转换通过 protobuf 的message_to_json()将事件转换为 JSON 字典可配置保留 proto 原始字段名或转为 camelCase。转换在ThreadPoolExecutor中执行避免阻塞事件循环async_publisher_client.py#L143-L156HTTP POST将过滤后的 JSON 数组发送到配置端点async_publisher_client.py#L172-L182aiohttp.ClientSession延迟创建、按需复用错误处理捕获异常并返回失败状态会话管理使用aiohttp.ClientSession管理 HTTP 会话close()时释放。批量发布事件按批次发布批次大小由max_batch_size限制MAX_EVENT_SEND_BATCH_SIZE默认 1,000批次由wait_for_batch()生成它最多等待一个超时窗口来凑齐事件更大的批次能减少 HTTP 请求开销但会增加发布延迟。超时窗口默认仅 0.1 秒意在兼顾吞吐与实时性。重试逻辑发布器实现了带指数退避的重试ray_event_publisher.py#L155-L211对失败的发布最多重试max_retries次。默认PUBLISHER_MAX_RETRIES -1即无限重试configs.py#L12重试间隔为指数退避叠加 jitter初始0.01秒、上限5.0秒、jitter 比例0.1指数上限封顶为2^30若重试次数耗尽配置为非负值时丢弃该批次事件并记录丢弃指标。配置汇总HTTP 发布相关配置通过环境变量完成全部定义在 aggregator_agent.py#L44-L82 与 configs.py环境变量默认值说明RAY_DASHBOARD_AGGREGATOR_AGENT_EVENTS_EXPORT_ADDR空外部 HTTP 服务端点 URL例如http://localhost:8080/eventsRAY_DASHBOARD_AGGREGATOR_AGENT_EXPOSABLE_EVENT_TYPES见DEFAULT_HTTP_EXPOSABLE_EVENT_TYPES允许暴露给外部服务的事件类型逗号分隔设为ALL则暴露全部类型ALL自 Ray 2.54.0 起支持RAY_DASHBOARD_AGGREGATOR_AGENT_PUBLISH_EVENTS_TO_EXTERNAL_HTTP_SERVICETrue是否发布事件到外部 HTTP 服务的开关RAY_DASHBOARD_AGGREGATOR_AGENT_MAX_EVENT_BUFFER_SIZE100000聚合 Agent 事件缓冲最大条数RAY_DASHBOARD_AGGREGATOR_AGENT_MAX_EVENT_SEND_BATCH_SIZE1000单批发送的最大事件数RAY_DASHBOARD_AGGREGATOR_AGENT_PRESERVE_PROTO_FIELD_NAMEFalse转 JSON 时是否保留 proto 原始字段名False表示转换为 camelCaseRAY_DASHBOARD_AGGREGATOR_AGENT_PUBLISHER_TIMEOUT_SECONDS3发布超时秒RAY_DASHBOARD_AGGREGATOR_AGENT_PUBLISHER_MAX_RETRIES-1最大重试次数小于 0 表示无限重试注意events_export_addr也可以通过 dashboard agent 的配置传入aggregator_agent.py#L132-L134环境变量是另一种途径。只有当PUBLISH_EVENTS_TO_EXTERNAL_HTTP_SERVICE为真且端点地址非空时HTTP 发布才会启用。创建新事件类型若需要为 Ray 引入新的可观测事件可按下述六个步骤操作完整示例可参考 Actor 定义事件的实现ray_actor_definition_event.h 与 ray_actor_definition_event.cc。第 1 步定义 Protobuf 消息在src/ray/protobuf/public/下新建.proto文件命名遵循events_name_event.proto约定参考 events_task_definition_event.proto。定义事件专属消息与所需字段syntax proto3; package ray.rpc.events; message MyNewEvent { // 在这里定义事件专属字段 string entity_id 1; // ... 其他字段 }第 2 步加入基础事件更新 events_base_event.proto为新的 proto 文件添加 import在EventType枚举中新增值例如MY_NEW_EVENT 18注意当前最大值为 17需顺延在RayEvent消息中新增嵌套字段例如MyNewEvent my_new_event 26。第 3 步实现 RayEventInterface创建实现RayEventInterface的 C 类最简单的方式是继承RayEventT模板类参考 ray_actor_definition_event.h#L26-L34。需要实现四个核心方法GetEntityId()返回实体唯一标识如任务 ID attempt、Actor IDMergeData()实现相同实体 ID 事件的合并逻辑——定义类事件合并时通常不变生命周期类事件则追加状态迁移SerializeData()将事件数据转换为RayEventprotobufGetEventType()返回本事件对应的EventType枚举值。完整的模板接口签名Merge、Serialize等见 ray_event_interface.h#L25-L71。注意Serialize()返回错误时例如嵌套事件负载解析失败记录器会跳过该事件序列化失败的事件会在导出时被跳过。第 4 步更新可暴露事件类型按需若新事件需要暴露给外部 HTTP 服务把它加入 configs.py#L34-L41 的DEFAULT_HTTP_EXPOSABLE_EVENT_TYPES否则用户可通过RAY_DASHBOARD_AGGREGATOR_AGENT_EXPOSABLE_EVENT_TYPES环境变量自行配置或设为ALL全部暴露。第 5 步让 RayEventRecorder 发布新事件使用RayEventRecorder::AddEvents()ray_event_recorder.cc#L65将新事件类型加入缓冲之后它会被周期导出。第 6 步让 AggregatorAgent 发布新事件AggregatorAgent中的发布配置aggregator_agent.py#L132-L192决定了事件发往哪些目的地外部 HTTP、GCS、dashboard head。外部 HTTP 端点由AsyncHttpPublisherClient统一处理事件一旦进入缓冲并属于可暴露类型即会被发布通常无需改动 Agent 主逻辑。小结Ray 的事件导出基础设施是一条「C 采集 → gRPC 传输 → Python 聚合 → HTTP 发布」的完整流水线C 侧RayEventRecorder负责线程安全的有界缓冲、按实体 ID 合并与周期导出Python 侧AggregatorAgent通过MultiConsumerEventBuffer实现多消费者、asyncio 安全的缓冲RayEventPublisher配合AsyncHttpPublisherClient完成过滤、JSON 转换、批量 POST 与指数退避重试。这套设计在保证事件可靠传递的同时通过合并与批量等手段显著降低了数据量与网络开销并且通过EventType枚举与RayEventInterface抽象为扩展新事件类型留出了清晰的路径。若要进一步实践可参考 ray-event-export.rst 中面向用户的导出配置指南或在 events_base_event.proto 与 aggregator_agent.py 中查看类型定义与发布开关的完整实现。【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表