ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

如何编写 SSE 客户端实时读取 OpenMetadata 摄取管道日志流

如何编写 SSE 客户端实时读取 OpenMetadata 摄取管道日志流 如何编写 SSE 客户端实时读取 OpenMetadata 摄取管道日志流【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata如果你在开发一个需要实时查看 OpenMetadata 摄取管道metadata、profiler、lineage、dbt 等运行日志的客户端传统的做法是轮询分页接口GET /logs/{id}/last?aftercursor选一个间隔、重复发请求、比对游标直到 run 结束。OpenMetadata 提供了一个 Server-Sent EventsSSE端点替代轮询GET /api/v1/services/ingestionPipelines/logs/{fqn}/stream/{runId}。响应类型为text/event-stream每帧是一个 JSON 格式的LogStreamEvent服务端每 25 秒发送一次 SSE 心跳注释: heartbeat防止代理断开空闲连接客户端解析时跳过即可。本文基于仓库内的 ingestion-log-streaming.md 和 streamable-logs.md 说明如何编写这样的客户端。前提条件确认服务端有日志后端流式端点对所有部署是同一个入口但服务端必须先配置了日志后端S3/MinIO 或 pipeline service。如果部署没有配置日志后端请求不会报 HTTP 错误而是直接在流上收到一个error事件然后关闭源文档给出的提示是 No log backend is configured on this deployment, so ingestion logs cannot be streamed.。使用 S3 对象存储时服务端配置位于openmetadata.yaml的pipelineServiceClientConfiguration.logStorageConfiguration下见 openmetadata.yamllogStorageConfiguration: type: ${PIPELINE_SERVICE_CLIENT_LOG_TYPE:-default} # Possible values are default, s3 enabled: ${PIPELINE_SERVICE_CLIENT_LOG_ENABLED:-false} # Enable it for pipelines deployed in the server # if type is s3, provide the following configuration bucketName: ${PIPELINE_SERVICE_CLIENT_LOG_BUCKET_NAME:-} prefix: ${PIPELINE_SERVICE_CLIENT_LOG_PREFIX:-} enableServerSideEncryption: ${PIPELINE_SERVICE_CLIENT_LOG_SSE_ENABLED:-false} sseAlgorithm: ${PIPELINE_SERVICE_CLIENT_LOG_SSE_ALGORITHM:-AES256} # Allowed values: AES256 or aws:kms awsConfig: enabled: ${PIPELINE_SERVICE_CLIENT_AWS_IAM_AUTH_ENABLED:-false} awsAccessKeyId: ${PIPELINE_SERVICE_CLIENT_LOG_AWS_ACCESS_KEY_ID:-} awsSecretAccessKey: ${PIPELINE_SERVICE_CLIENT_LOG_AWS_SECRET_ACCESS_KEY:-} awsRegion: ${PIPELINE_SERVICE_CLIENT_LOG_AWS_REGION:-} endPointURL: ${PIPELINE_SERVICE_CLIENT_LOG_AWS_ENDPOINT_URL:-} # port forward localhost:9000 for minio仓库还提供了带环境变量示例的完整配置 openmetadata-s3-logs.yaml其中type: s3时通过bucketName、region、prefix和awsConfig指定存储位置。日志写入端的设计连接器如何批量 POST 日志、/close如何收尾、abandoned-run sweeper 如何回收见 streamable-logs.md客户端只需要知道读取侧行为。拿到两个路径参数后就能请求端点{fqn}管道的 fullyQualifiedName 或 IdUUID{runId}要跟踪的 run。日志在对象存储中时是 UUID否则是 pipeline service 自己的 run 标识如 Airflow 的scheduled__…。要 tail 最新一次运行先从管道的pipelineStatuses字段读出其runId—— 这是分页接口/logs/{id}/last内部使用的同一个字段。先用 curl 验证端点在写客户端之前用 curl 确认端点可用-N关闭 curl 缓冲这样实时 tail 才可见curl -N -H Authorization: Bearer $OM_TOKEN \ http://localhost:8585/api/v1/services/ingestionPipelines/logs/my.pipeline.fqn/stream/$RUN_ID其中$OM_TOKEN换成你的 Bearer tokenmy.pipeline.fqn换成实际管道 FQN或 UUID$RUN_ID换成从pipelineStatuses读到的 run 标识。若 FQN 不对端点返回 404其余在管道解析成功之后才发生的问题无日志后端、服务端流容量已满、单客户端连接数超限都通过流上的error事件报告而不是 HTTP 状态码——所以客户端只需要一条错误处理路径。文档给出的事件示例runId中…为省略值非固定输出// eventType: logs — new content {eventType:logs,runId:a1b2…,logs:[2026-08-10 …] INFO Ingesting table x,after:4211,replay:false,truncated:false} // eventType: complete — the server is closing the stream {eventType:complete,runId:a1b2…,after:4680,reason:runFinished} // eventType: error — the stream cannot be served; it is closed right after {eventType:error,runId:a1b2…,message:The server is already streaming the maximum number of pipeline runs. …}事件字段客户端要解析的全部内容事件 schema 定义在 logStreamEvent.json。所有帧都挂在未命名的 SSE 事件上因此浏览器里普通的EventSource.onmessage就能收到全部帧。字段含义eventTypelogs、complete或errorrunId内容所属的 runlogs自上一事件以来追加的内容complete/error上不存在after指向logs之后的游标。必须保存重连时作为?after传回replaytrue表示该块来自服务端 replay 缓冲连接时流已在运行truncated首个事件上为true表示服务端无法精确算出你缺了什么。按重置处理清空查看器、渲染随后 replay 的内容并从GET /logs/{id}/last或 download 端点回填更早历史reason流结束原因见下表messageerror事件的可读细节以及提前结束的complete的说明after游标是不透明的对象存储后端下它是行偏移Airflow 后端下是块偏移两者不可互换永远不要手工构造。编写浏览器客户端EventSource不能设置Authorization头所以用fetch加流式 reader。以下是文档给出的实现其中getBasePath、getEncodedFqn、getOidcToken、appendToViewer、onStreamEnd、onStreamError是你需要按自己应用提供/替换的辅助函数分别负责 API 基础路径、FQN 编码FQN 含/时需 URL 编码、获取 Bearer token、向查看器追加日志、以及处理流结束/出错const controller new AbortController(); let cursor: string | undefined; const tail async (fqn: string, runId: string) { const url new URL( ${getBasePath()}/api/v1/services/ingestionPipelines/logs/${getEncodedFqn( fqn )}/stream/${encodeURIComponent(runId)}, window.location.origin ); if (cursor) { url.searchParams.set(after, cursor); } const response await fetch(url, { headers: { Authorization: Bearer ${await getOidcToken()} }, signal: controller.signal, }); const reader response.body!.getReader(); const decoder new TextDecoder(); let buffer ; for (;;) { const { done, value } await reader.read(); if (done) { break; } buffer decoder.decode(value, { stream: true }); const frames buffer.split(\n); buffer frames.pop() ?? ; for (const frame of frames) { if (!frame.startsWith(data: )) { continue; // heartbeat comment or blank separator } const event JSON.parse(frame.slice(6)); cursor event.after ?? cursor; if (event.eventType logs) { appendToViewer(event.logs); } else if (event.eventType complete) { onStreamEnd(event.reason); // reconnect here for a non-runFinished reason } else { onStreamError(event.message); } } } };解析逻辑的三个关键点按行拆帧跳过非data:开头的行。心跳注释: heartbeat和空分隔行都会出现直接跳过。每收到一帧就更新cursorevent.after。它是断线后唯一的续传凭据。truncated: true时重置查看器而不是追加。当你的游标比共享 reader 的 replay 缓冲更旧或来自负载均衡后另一台服务器服务端无法判断中间缺了什么会带着truncated: true重放它已有的部分——这是唯一必须清空查看器的场景。处理流结束与重连complete事件的reason决定了客户端下一步做什么reason发生了什么客户端应做什么runFinishedrun 到达终态且日志已静默或服务端没有该 run 的状态行且静默了一分钟什么都不做日志已完整idleTimeout5 分钟无新内容且 run 未报告终态若仍关心该 run带?after重连maxDuration流达到 1 小时生命周期上限带?after重连maxBytes流已交付 32 MB剩余内容改用GET /logs/{id}/last/download流体没有complete事件就关闭说明被截断客户端停止消费并越过自身积压上限、服务端消失或网络中断。处理方式同样是带?after最后游标重连。重连之间要退避。上述原因往往持续存在追不上的查看器下次还会再越过积压上限每次重连都会重新拉取 replay 积压。立即循环重连会把一个卡顿的客户端变成压力源。文档建议使用带上限的指数退避连续失败几次后放弃而不是无限重试。限制与已知边界每个 (storage backend, pipeline FQN, run) 只有一个reader同一次运行在十个标签页打开会产生十条 SSE 连接但对 S3/Airflow 仍只有一个 reader最后一个查看器断开时 reader 停止所以没人看的 run 不会被读取。服务端各项上限LogStreamSettings强制单次 tick 推送 1 MB、单条流总量 32 MB、单条流 1 小时、空闲 300 秒、每 run replay 缓冲 256 KB、单客户端积压 4 MB、每服务端最多 200 个并发 tail 的 run、500 条连接。超过maxActiveRuns或maxActiveConnections的请求会收到error事件并关闭而不是排队。多节点部署下 tailer 是每服务端一份负载均衡后两台服务器都有人看同一个 run 时各持一个 reader这是设计内的取舍读取无状态、从 S3 读partial.txt任意实例都可以读路径不需要 sticky session写路径才需要见 streamable-logs.md。验证与延伸阅读验证方式与前面 curl 一节一致能连续收到logs帧、run 结束后收到reason: runFinished的complete帧说明客户端解析和游标逻辑正确。想核对完整行为可参考仓库中的端到端测试 IngestionPipelineLogStreamIT.java 和流式引擎实现 logstorage/stream/。若客户端侧只是要读取已结束 run 的历史分页端点GET /logs/{fqn}/{runId}和 download 端点仍按 streamable-logs.md 的 Read Paths 表工作SSE 不是唯一选择。【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表