ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Lago Ingest Connectors 实战指南:基于 SQS、Kinesis 与 HTTP 的高吞吐事件采集管道

Lago Ingest Connectors 实战指南:基于 SQS、Kinesis 与 HTTP 的高吞吐事件采集管道 Lago Ingest Connectors 实战指南基于 SQS、Kinesis 与 HTTP 的高吞吐事件采集管道【免费下载链接】lagoOpen Source Metering and Usage Based Billing API ⭐️ Consumption tracking, Subscription management, Pricing iterations, Payment orchestration Revenue analytics项目地址: https://gitcode.com/GitHub_Trending/la/lagoLago 是一套开源的使用量计量Metering与基于用量计费Usage Based Billing平台其核心数据来源是业务侧上报的 usage event。本文围绕仓库 connectors/README.md 展开系统讲解 Lago 的事件格式规范、三种官方采集入口SQS、Kinesis、HTTP Server的环境变量配置并结合仓库内的 sqs.yml、kinesis.yml、http.yml 三个完整配置以及 events-processor 的消费端源码说明事件从业务侧产生到进入 Lago 计费核心的完整链路。读完本文你将掌握如何按官方规范构造事件、如何部署并配置三种 connector、如何理解其中的 DLQ、批量发送与动态定价字段并能在自己的环境里完成验证与排错。Lago 事件采集架构与 Connectors 的定位在 Lago 的整体架构中事件从产生到最终影响账单会经过典型的采集 → 传输 → 处理三个阶段采集Ingest业务系统通过 Lago API 或本仓库提供的 connector 上报事件。connectors 目录下提供了三种独立的采集入口——AWS SQS、AWS Kinesis、HTTP Server它们将外部事件源接入统一的 RedpandaKafka 协议消息总线传输Transportconnector 通过kafka_franz输出组件将标准化后的事件写入 Redpanda 的 raw events topic即 events-processor 消费的LAGO_KAFKA_RAW_EVENTS_TOPIC处理Processingevents-processor 作为高吞吐事件后处理服务从该 topic 消费事件完成 enrichment补充计费指标、订阅等信息后写入 enriched 系列 topic最终驱动计费、发票等下游逻辑。在 docs/architecture.md 的 Worker 参考表中这两类角色分别被称为Events Consumer Worker消费外部队列如 Kafka/SQS 的事件与Events Processor Worker处理并聚合 usage events它们共同构成 Lago 的事件流式处理管道。connectors 正是采集阶段的关键拼图它让 Lago 可以脱离同步 API 调用以异步、可扩展的方式接收海量计费事件。事件格式规范与 Lago API 完全一致的 Event 结构connectors 的首要约束是所有事件都必须遵循与 Lago API 相同的事件格式参见 README 的 Events format 一节。一个合法的 usage event 示例{ event: { external_subscription_id: string, transaction_id: unique_transaction_identifier, code: billable_metric_code, // Event Unix Timestamp timestamp: 1620000000, properties: { // Should respect the format of your billable metric my_property: my_value }, // Optional, used for the dynamic pricing feature, defaulted to 0 precise_total_amount_cents: 1000 } }各字段含义如下字段类型必填说明external_subscription_idstring是业务侧订阅的外部标识用于将事件关联到具体订阅transaction_idstring是事件唯一事务标识Lago 依赖它做幂等去重codestring是计费指标billable metric的 code与你在 Lago 后台定义的指标对应timestampinteger是事件发生的 Unix 时间戳秒propertiesobject是事件属性其键值必须符合该计费指标定义的属性格式如聚合维度、数值字段precise_total_amount_centsinteger否动态定价dynamic pricing专用表示事件对应的精确金额分缺省为 0说明该字段与 Lago 动态定价功能相关。当计费指标需要按事件粒度携带精确金额时通过该字段传递不传则按 0 处理。该结构与 events-processor 的 Event 模型 严格对应OrganizationID、ExternalSubscriptionID、TransactionID、Code、Properties、PreciseTotalAmountCents、Timestamp、IngestedAt等字段都有对应的 JSON tag。也就是说connector 侧映射出来的字段名与 Go 侧反序列化所需的 JSON 字段名是同一套契约任何字段名偏差都会导致事件在消费端解析失败并被送入死信队列。值得注意的是事件经过 connector 映射后会额外带上organization_id从环境变量注入和ingested_at采集时间这两个字段不在原始上报格式中而是由 connector 在管道内自动补充详见下文映射配置。通用配置与镜像基础所有 connector 共享一项环境变量配置环境变量说明必填LOG_LEVELconnector 的日志级别默认info否三个 connector 均构建于同一个镜像之上。connectors/Dockerfile 显示其基础镜像是docker.io/redpandadata/connect:4.83.0即 Redpanda Connect原 Bento流处理引擎并把目录下的三个*.yml配置直接拷入镜像根目录FROM docker.io/redpandadata/connect:4.83.0 COPY *.yml ./这意味着每个 connector 的 YAML 都是标准 Redpanda Connect 配置文件遵循其input → pipeline → output三段式结构同时可以使用其处理器、资源复用、测试框架tests与 Prometheus 指标暴露能力。三类 connector 的配置都遵循同一模板不同的输入源AWS SQS / AWS Kinesis / HTTP Server 统一的lago_event_mapping映射处理器 统一的 RedpandaKafka输出。SQS Connector从 AWS 队列消费事件SQS connector 从 AWS Simple Queue Service 拉取事件消息是最常见的异步接入方式。其环境变量如下环境变量说明必填SQS_ENDPOINTSQS 服务地址队列 URL是SQS_REGIONSQS 服务所在区域是SQS_KEY_IDAWS 访问密钥 ID是SQS_KEY_SECRETAWS 访问密钥 Secret是SQS_DLQ_ENDPOINT死信队列DLQ服务地址否ORGANIZATION_IDLago 组织 ID是KAFKA_BROKERSRedpanda Broker 地址是KAFKA_USERRedpanda 用户名是KAFKA_PASSWORDRedpanda 密码是KAFKA_TOPIC发送事件的 Redpanda Topic是完整的输入配置见 sqs.ymlinput: label: sqs_lago_ingestor aws_sqs: url: ${SQS_ENDPOINT} max_number_of_messages: 10 region: ${SQS_REGION} credentials: id: ${SQS_KEY_ID} secret: ${SQS_KEY_SECRET}其中max_number_of_messages: 10是单次拉取的最大消息数SQS 单次批量接收的上限保证每轮可以并行处理多条事件提升吞吐。事件映射处理器三个 connector 的核心逻辑都集中在名为lago_event_mapping的mapping处理器Bloblang 语言中。SQS 场景下原始消息结构是{event: {...}}外层包了一层event因此映射从this.event取值processor_resources: - label: lago_event_mapping mapping: | root.organization_id ${ORGANIZATION_ID} root.external_subscription_id this.event.external_subscription_id root.transaction_id this.event.transaction_id root.code this.event.code root.timestamp this.event.timestamp root.properties this.event.properties root.ingested_at timestamp_unix() root.precise_total_amount_cents if this.event.precise_total_amount_cents.type() number { this.event.precise_total_amount_cents } else { 0 }这段映射做了三件事注入组织标识root.organization_id直接取自环境变量ORGANIZATION_ID消费者无需在每条事件中重复携带字段平铺把event内嵌对象中的字段提升到顶层并新增ingested_at timestamp_unix()记录采集时刻动态定价字段兜底precise_total_amount_cents用 Bloblang 条件表达式处理——若原始事件中该字段是数字则原样保留否则统一填充字符串0注意输出端该字段为字符串类型与 Event 模型 中PreciseTotalAmountCents string的类型一致。输出与死信队列路由SQS connector 的输出使用switch条件路由sqs.ymloutput: switch: cases: - check: errored() env(SQS_DLQ_ENDPOINT).type() string output: resource: lago_sqs_ingestor_dlq - output: resource: lago_redpanda_events_raw当管道处理出错errored()并且配置了SQS_DLQ_ENDPOINT时消息被投递到 SQS 死信队列lago_sqs_ingestor_dlq避免坏消息阻塞主链路否则写入 Redpanda 的 raw events topiclago_redpanda_events_raw。两个 output 资源的定义sqs.yml也值得注意。写 Redpanda 时output_resources: - label: lago_redpanda_events_raw kafka_franz: seed_brokers: [${KAFKA_BROKERS}] topic: ${KAFKA_TOPIC} key: ${ORGANIZATION_ID}-${! json(external_subscription_id) } tls: enabled: true sasl: - mechanism: SCRAM-SHA-512 username: ${KAFKA_USER} password: ${KAFKA_PASSWORD}消息 key采用组织ID-订阅外部ID的组合这保证同一订阅的事件在 Kafka 分区上保持有序同一 key 路由到同一分区对下游计费聚合的时序一致性至关重要安全SQS 连接器默认启用 TLS并使用 SCRAM-SHA-512 机制认证。注意这里tls.enabled: true是硬编码的与 Kinesis/HTTP 连接器由KAFKA_TLS控制不同DLQlago_sqs_ingestor_dlq复用SQS_REGION/SQS_KEY_ID/SQS_KEY_SECRET凭证向SQS_DLQ_ENDPOINT写回死信。内置测试用例sqs.yml 还附带了两组 Redpanda Connect 内置测试可以用bento test或对应 connect 的测试命令直接运行用于验证映射逻辑全字段用例输入包含全部字段含precise_total_amount_cents: 100的事件断言输出 JSON 中organization_id被注入为环境变量123其余字段原样透传precise_total_amount_cents保持数值100缺省字段用例输入不含precise_total_amount_cents的事件断言输出中该字段被兜底为字符串0。这两个用例正是对动态定价字段可选、缺省为 0这一规范的直接验证也是改动映射后回归测试的最小单元。Kinesis Connector从 Kinesis 流消费事件Kinesis connector 面向 AWS Kinesis Data Streams 场景适用于已经用 Kinesis 汇集大量事件的架构。其环境变量如下环境变量说明必填KINESIS_STREAM要消费的 Kinesis 流名称是AWS_REGIONKinesis 与 DynamoDB 所在的 AWS 区域是AWS_ROLE访问 Kinesis 时扮演的 IAM 角色 ARN是AWS_ROLE_EXTERNAL_ID被扮演角色的外部 ID是DYNAMODB_TABLE用于 checkpoint 的 DynamoDB 表名是KAFKA_BROKERSRedpanda Broker 地址是KAFKA_USERRedpanda 用户名是KAFKA_PASSWORDRedpanda 密码是KAFKA_TOPIC发送事件的 Redpanda Topic是KAFKA_TLS是否对 Kafka 连接启用 TLS默认false否KAFKA_BATCH_COUNT批量发送前的消息条数阈值默认100否KAFKA_BATCH_BYTE_SIZE批量发送的字节数上限默认10000001MB否KAFKA_BATCH_PERIOD批量发送的时间窗口默认1s否输入配置见 kinesis.ymlinput: aws_kinesis: streams: [${KINESIS_STREAM}] region: ${AWS_REGION} dynamodb: table: ${DYNAMODB_TABLE} create: true billing_mode: PAY_PER_REQUEST region: ${AWS_REGION} checkpoint_limit: 1024 auto_replay_nacks: true commit_period: 1s steal_grace_period: 1s start_from_oldest: true credentials: role: ${AWS_ROLE} role_external_id: ${AWS_ROLE_EXTERNAL_ID}Kinesis 是至少一次at-least-once投递语义的流需要本地维护消费位点。该配置的几个关键点checkpoint 存储使用 DynamoDB 表DYNAMODB_TABLE记录每个 shard 的消费位点create: true表示表不存在时自动创建billing_mode: PAY_PER_REQUEST按量计费消费语义start_from_oldest: true从最老数据开始消费适合首次启动auto_replay_nacks: true对处理失败的消息自动重放commit_period: 1s与steal_grace_period: 1s控制 checkpoint 提交频率与分片再平衡的宽限期凭证通过AWS_ROLEIAM 角色 ARNAWS_ROLE_EXTERNAL_ID以角色扮演assume role方式访问 Kinesis而不是直接使用静态密钥适合跨账号或最小权限场景。Kinesis 场景下由于原始消息已经是被展平的事件无外层event包裹其映射处理器直接从顶层取值kinesis.ymlprocessor_resources: - label: lago_event_mapping mapping: | root.organization_id ${ORGANIZATION_ID} root.external_subscription_id this.external_subscription_id root.transaction_id this.transaction_id root.code this.code root.timestamp this.timestamp root.properties this.properties root.ingested_at timestamp_unix() root.precise_total_amount_cents if this.precise_total_amount_cents.type() number { this.precise_total_amount_cents } else { 0 }与 SQS 版本唯一的差异是取值路径少了event.前缀——这也提醒使用者不同入口对原始消息的封装层级不同SQS 外层包eventKinesis 不包接入前必须先确认上游投递的实际 JSON 结构。HTTP Server Connector以 HTTP 端点接收事件HTTP connector 在容器内启动一个 HTTP 服务业务侧可以直接向该端点 POST 事件适合不依赖 AWS、希望以最简单方式接入的场景。其环境变量如下环境变量说明必填KAFKA_BROKERSRedpanda Broker 地址是KAFKA_USERRedpanda 用户名是KAFKA_PASSWORDRedpanda 密码是KAFKA_TOPIC发送事件的 Redpanda Topic是KAFKA_TLS是否对 Kafka 连接启用 TLS默认false否KAFKA_BATCH_COUNT批量发送前的消息条数阈值默认100否KAFKA_BATCH_BYTE_SIZE批量发送的字节数上限默认10000001MB否KAFKA_BATCH_PERIOD批量发送的时间窗口默认1s否输入配置见 http.ymlinput: label: http_lago_ingestor http_server: address: 0.0.0.0:3000 path: /events allowed_verbs: - POST即监听0.0.0.0:3000的/events路径仅接受 POST 请求。业务侧可以直接向该端点上报事件请求体使用 README 中定义的标准事件 JSON。HTTP 版本的输出最有特点http.ymloutput: broker: pattern: fan_out_sequential # Process sequentially to ensure Kafka write happens first outputs: - resource: lago_kafka_events_raw - processors: - mapping: root # Return empty string sync_response: {} # Immediately respond to HTTP requestpattern: fan_out_sequential表示两个输出按顺序依次执行先保证 Kafka 写入完成再返回 HTTP 响应——避免响应成功但事件丢失的假成功第二个输出通过mapping: root 把响应体置空并配合sync_response: {}立即对 HTTP 请求返回空响应200让调用方感知写入结果。Redpanda 输出共性TLS、SCRAM 认证与批量发送SQS、Kinesis、HTTP 三种连接器的输出端最终都汇聚到同一个 RedpandaKafka 协议写入配置以 Kinesis 版本为例kinesis.ymloutput_resources: - label: lago_kafka_events_raw kafka_franz: seed_brokers: [${KAFKA_BROKERS}] topic: ${KAFKA_TOPIC} key: ${! json(organization_id) }-${! json(external_subscription_id) } tls: enabled: ${KAFKA_TLS:false} sasl: - mechanism: SCRAM-SHA-512 username: ${KAFKA_USER} password: ${KAFKA_PASSWORD} batching: count: ${KAFKA_BATCH_COUNT:100} byte_size: ${KAFKA_BATCH_BYTE_SIZE:1000000} # 1MB default period: ${KAFKA_BATCH_PERIOD:1s}这里集中体现了 Kinesis/HTTP 版本的批量发送参数三者任一达到即触发发送参数默认值说明KAFKA_BATCH_COUNT100攒够多少条消息触发一次批量发送KAFKA_BATCH_BYTE_SIZE10000001MB批量消息总字节数上限KAFKA_BATCH_PERIOD1s最长等待窗口低流量时也不至于无限积压批量发送能显著降低与 Redpanda 之间的网络往返次数是吞吐优化的重要手段消息 key 与 SQS 版本一致组织ID-订阅外部ID保证订阅内事件分区有序。与消费端的认证契约一致connector 写入端使用 SCRAM-SHA-512而 events-processor 消费端同样支持 SCRAM 认证。kafka.go 中定义了SCRAM-SHA-256与SCRAM-SHA-512两种算法由LAGO_KAFKA_SCRAM_ALGORITHM指定TLS 通过LAGO_KAFKA_TLS开启。因此在搭建 Redpanda 时务必保证 connector 与 events-processor 两侧的认证与加密配置对齐否则会出现写入成功但消费端连不上的隐性故障。事件进入 Lago 后的处理链路理解 connector 的终点有助于确认它是否把事件送对了地方。events-processor 的入口 main_processor.go 会读取以下关键环境变量建立 Kafka 连接LAGO_KAFKA_BOOTSTRAP_SERVERSBroker 列表如redpanda:9092,kafka:9092LAGO_KAFKA_RAW_EVENTS_TOPIC消费的原始事件 topic——这正是 connector 的KAFKA_TOPIC应当指向的位置LAGO_KAFKA_CONSUMER_GROUP消费组名称LAGO_KAFKA_ENRICHED_EVENTS_TOPIC、LAGO_KAFKA_ENRICHED_EVENTS_EXPANDED_TOPIC、LAGO_KAFKA_EVENTS_CHARGED_IN_ADVANCE_TOPIC、LAGO_KAFKA_EVENTS_DEAD_LETTER_TOPIC处理后的输出 topic 与死信 topic。更完整的参数表见 events-processor/README.md。消费端的可靠性设计也值得与 connector 的 DLQ 机制对照理解。在 consumer.go 中消费组按分区并行处理PartitionConsumer只有成功处理的记录才会被 commit若批次中有记录处理失败则只提交到最大可提交记录为止失败记录会在 rebalance 后重新拉取。而 processor.go 进一步规定可重试的错误且事件未超过12 小时不提交 offset等待重新消费超过 12 小时或不可重试的错误写入LAGO_KAFKA_EVENTS_DEAD_LETTER_TOPIC死信 topic并携带error_message、error_code、failed_at等诊断信息见 event_producer_service.go 与 event.go 中的FailedEvent结构。因此整条链路形成了双层容错connector 层对映射失败的事件回写 SQS DLQevents-processor 层对 enrichment 失败的事件送入 Kafka DLQ topic。这也是 README 强调事件格式必须与 Lago API 一致的根本原因——格式不合规的事件会在两层 DLQ 中沉淀下来最终需要人工或脚本清洗。实操构建、配置与验证1. 构建 connector 镜像connectors 目录是自包含的直接基于 Dockerfile 构建即可docker build -t lago-ingest-connector .构建产物包含三个 YAML 配置运行哪一个 connector 取决于传入的配置文件名例如bento -c sqs.yml或通过镜像入口指定。2. 准备 Redpanda 与 Topic按 scripts/create-topics.sh 的做法topic 需要预先创建该脚本用rpk实现不存在才创建的幂等逻辑。至少需要准备raw events topicconnector 写入、events-processor 消费以及 enriched、enriched_expanded、charged_in_advance、dead_letter 等下游 topic。3. 以 SQS connector 为例启动配置好 sqs.yml 所需的全部环境变量后启动docker run --rm \ -e SQS_ENDPOINThttps://sqs.us-east-1.amazonaws.com/123456789012/lago-events \ -e SQS_REGIONus-east-1 \ -e SQS_KEY_IDAKIA... \ -e SQS_KEY_SECRET... \ -e ORGANIZATION_IDyour-org-id \ -e KAFKA_BROKERSredpanda:9092 \ -e KAFKA_USERlago \ -e KAFKA_PASSWORDsecret \ -e KAFKA_TOPICevents_raw \ -e LOG_LEVELinfo \ lago-ingest-connector sqs.yml各环境变量的默认值与取值范围以上文表格为准LOG_LEVEL默认infoKinesis/HTTP 版本中KAFKA_TLS默认false、KAFKA_BATCH_COUNT默认 100、KAFKA_BATCH_BYTE_SIZE默认 1MB、KAFKA_BATCH_PERIOD默认 1s而 SQS 版本的 Redpanda 输出默认开启 TLS配置时需与 Broker 实际配置一致。4. 向 HTTP connector 发送测试事件若使用 HTTP 入口直接 POST 标准事件格式即可curl -X POST http://localhost:3000/events \ -H Content-Type: application/json \ -d { event: { external_subscription_id: sub_123, transaction_id: txn_456, code: api_calls, timestamp: 1620000000, properties: {region: eu-west-1} } }由于fan_out_sequential保证先写 Kafka 再响应HTTP 返回 200 即代表事件已成功落入 Redpanda raw topic。5. 验证映射逻辑SQS 配置自带的tests段可以直接运行验证映射正确性两个用例分别覆盖全字段与缺省precise_total_amount_cents两种输入。对于自己调整过的映射可以照葫芦画瓢补充用例确保字段平铺、组织 ID 注入与动态定价兜底行为不被破坏。6. 监控三个 YAML 均开启了metrics.prometheus: {}connector 会暴露 Prometheus 指标端点可接入现有监控体系观察拉取/发送速率、错误率与批量发送情况与 docs/monitoring.md 中描述的 Lago 监控实践配套使用。小结Lago 的 connectors 用一套标准化的 Redpanda Connect 配置把三种截然不同的外部事件源AWS SQS、AWS Kinesis、HTTP统一接入同一条事件管道统一的lago_event_mapping完成字段标准化与组织标识注入统一的kafka_franz输出完成安全认证与批量发送加上各自的 DLQ/checkpoint 机制保证可靠性。理解这套设计的关键在于三张契约表README 定义的事件 JSON 格式、Event 模型 的字段定义以及 events-processor 消费端的环境变量约定。只要事件格式合规、两侧 Kafka 配置对齐、topic 就位即可把任意业务侧事件稳定地送入 Lago 计费核心。【免费下载链接】lagoOpen Source Metering and Usage Based Billing API ⭐️ Consumption tracking, Subscription management, Pricing iterations, Payment orchestration Revenue analytics项目地址: https://gitcode.com/GitHub_Trending/la/lago创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表