我的文档库

BI 埋点收集服务:基于 Cloudflare Workers 的数据采集与查询架构

BI 埋点收集服务是一套部署在 Cloudflare Workers 上的前端数据采集与查询服务,负责接收客户端埋点、异步持久化事件数据,并为管理后台提供报表查询能力。

BI 埋点收集服务是一套部署在 Cloudflare Workers 上的前端数据采集与查询服务,负责接收客户端埋点、异步持久化事件数据,并为管理后台提供报表查询能力。


一、服务定位

BI 埋点服务主要承担三个职责:

客户端事件接收

异步数据持久化

报表数据查询

完整数据链路:

Browser SDK

     │ Encrypted Events

┌───────────────┐
│  Collect API  │
└───────┬───────┘

        │ Validate / Decrypt

┌───────────────┐
│ Cloudflare    │
│    Queue      │
└───────┬───────┘

        │ Consumer

┌───────────────┐
│ Data Storage  │
├───────────────┤
│ R2            │
│ Analytics     │
│ Engine        │
└───────┬───────┘


┌───────────────┐
│  Report API   │
└───────┬───────┘


     Admin / BI

二、为什么独立建设 BI 服务

前端应用产生的埋点数据具有几个明显特点:

  • 请求量大
  • 单条数据价值低
  • 不要求同步完成持久化
  • 数据结构相对稳定
  • 写入远多于查询
  • 查询具有明显的时间范围特征

因此没有必要让业务应用直接承担整个数据链路。

采用独立 BI 服务后:

Business Application

        │ Event

    BI Service


Async Data Pipeline

业务系统与数据基础设施实现解耦。


三、系统分层

结合当前服务实现,整个 BI 服务划分为七个层次:

┌──────────────────────────────────┐
│            接入层                │
│                                  │
│          Fetch Handler            │
└────────────────┬─────────────────┘

┌────────────────▼─────────────────┐
│            认证层                │
│                                  │
│     API Key / HMAC-SHA256       │
└────────────────┬─────────────────┘

┌────────────────▼─────────────────┐
│           数据采集层             │
│                                  │
│       Decrypt / Validate         │
│       Geo Enrichment             │
└────────────────┬─────────────────┘

┌────────────────▼─────────────────┐
│           异步处理层             │
│                                  │
│          Queue Consumer          │
└────────────────┬─────────────────┘

┌────────────────▼─────────────────┐
│           数据存储层             │
│                                  │
│         R2 / Analytics Engine    │
└────────────────┬─────────────────┘

        ┌────────┴────────┐
        ▼                 ▼
┌───────────────┐  ┌───────────────┐
│   数据维护层  │  │    查询层     │
│               │  │               │
│   Compaction  │  │    Report     │
└───────────────┘  └───────────────┘

这七层分别对应当前服务中的实际职责,而不是为了架构形式额外增加抽象层。


四、第一层:接入层

接入层由 Fetch Handler 负责。

它是 Worker 的 HTTP 入口:

HTTP Request


Fetch Handler

      ├── POST /collect

      └── GET /report

Fetch Handler 本身不负责具体的数据处理,而是根据请求路径将请求分发给对应模块。

src/index.ts

主要职责:

Request

Routing

Handler

Response

五、Collect API

埋点上报使用 POST 请求:

POST /service_stats/v1/collect

主要用于接收客户端加密事件。

完整流程:

Client


POST /collect


Authentication


Decrypt


Parse Events


Validate


Geo Enrichment


Queue

单次请求限制最多 100 条事件。


六、Report API

报表查询使用:

GET /service_stats/v1/report

查询流程:

GET /report


Authentication


Parse Query


Resolve Files


Read R2


Filter Events


Limit


Response

当前接口最大支持 100,000 条数据返回。


七、第二层:认证层

认证层负责验证请求是否合法。

系统目前支持两种认证模式:

模式Header使用场景
API KeyX-API-Key管理后台报表查询
HMAC-SHA256X-Signature + X-TimestampSDK 埋点上报

因此两条请求链路分别为:

Report


API Key


Report Handler

以及:

Collect


Timestamp


HMAC


Decrypt


Collect Handler

八、HMAC 签名

SDK 上报使用 HMAC-SHA256 对加密后的数据进行签名。

逻辑关系:

Payload


AES Encryption


Ciphertext


HMAC-SHA256


Signature

服务端收到请求后:

Request

   ├── Timestamp

   ├── Signature

   └── Ciphertext


      Verify HMAC


       Decrypt

当前方案中 AES 密钥由 SIGN_SECRET 派生,签名使用 HMAC-SHA256。


九、防重放

签名本身只能证明数据完整性和签名有效性,不能阻止同一个合法请求被重复发送。

因此增加 Timestamp:

X-Timestamp

服务端限制:

|now - timestamp| <= 5 minutes

超过有效窗口直接拒绝。

Request


Timestamp Check

   ├── Expired → Reject

   └── Valid


      HMAC Verify

这样可以降低历史请求被重复利用的风险。


十、第三层:数据采集层

数据采集层是整个服务的核心入口。

主要负责:

Decrypt
Parse
Validate
Enrich
Serialize
Enqueue

流程:

Encrypted Payload


AES-256-GCM Decrypt


JSON Parse


Event Validation


Geo Enrichment


Serialize


Queue

十一、事件解密

客户端发送的是加密后的事件数据。

服务端首先进行:

AES-256-GCM

解密。

当前数据格式中 nonce 位于密文前 12 字节:

┌──────────────┬──────────────────────┐
│    Nonce     │      Ciphertext      │
│   12 bytes   │                      │
└──────────────┴──────────────────────┘

服务端:

Payload

   ├── Nonce

   └── Ciphertext


       Decrypt


       JSON[]

十二、事件校验

解密后得到事件数组。

服务端需要限制单次请求规模:

Event[]


Validate

   ├── Invalid → Reject

   └── Valid

当前单个请求最多允许 100 条事件。

这样可以避免单个请求携带过大的 Payload。


十三、Geo 信息注入

Cloudflare Worker 可以从请求上下文获取地理信息。

因此在数据进入异步队列之前进行统一注入:

Event


Cloudflare Request Context

 ├── IP
 ├── ASN
 ├── Country
 ├── City
 ├── Region
 ├── Region Code
 ├── Timezone
 └── Postal Code

最终:

Client Event
     +
Cloudflare Metadata

Enriched Event

这样客户端 SDK 不需要自行获取和维护这些信息。


十四、第四层:异步处理层

数据采集完成后,并不直接同步写入 R2。

而是进入:

Cloudflare Queue

整体:

Collect Worker


    Queue


Queue Consumer

这样:

HTTP Request

与:

Data Persistence

实现解耦。


十五、为什么使用 Queue

如果同步写入:

Request


Decrypt


R2


Analytics Engine


Response

那么存储速度会直接影响用户请求。

采用 Queue:

Request


Validate


Queue


Response

后台:

Queue


Consumer

  ├── R2

  └── Analytics Engine

这样可以将请求延迟与数据持久化解耦。


十六、Queue Consumer

Consumer 负责真正的数据持久化:

Queue Message


Idempotency Check


Write R2


Write Analytics Engine

当前 Queue Consumer 的批处理配置:

max_batch_size = 100
max_batch_timeout = 30s

十七、Queue 幂等

Queue 消费必须考虑重复消息。

例如:

Message A

   ├── Consumer 1

   └── Retry


      Consumer 2

如果没有幂等控制:

A
A

可能被重复写入。

当前方案使用:

messageId

作为幂等依据。

写入前检查:

R2.head(key)

如果对象已经存在:

Skip

否则:

Write

十八、第五层:数据存储层

Consumer 将数据同时写入:

┌───────────────┐
│      Event    │
└───────┬───────┘

   ┌────┴────┐
   ▼         ▼
  R2    Analytics Engine

两个存储承担不同职责。

R2

用于保存原始事件数据:

Raw Events
Historical Data

Analytics Engine

用于:

Analytics
Metrics
Reporting

这样形成:

R2
→ Raw / Historical

Analytics Engine
→ Analytics / Query

十九、R2 数据分区

原始数据按照时间进行目录划分:

raw/
└── YYYY-MM-DD/
    └── HH/
        ├── message-001.jsonl
        ├── message-002.jsonl
        └── ...

例如:

raw/
└── 2026-09-16/
    └── 10/
        ├── abc.jsonl
        ├── def.jsonl
        └── ghi.jsonl

这样查询时可以根据时间范围快速定位数据。

当前方案使用:

raw/YYYY-MM-DD/HH/{messageId}.jsonl

作为对象 Key。


二十、Analytics Engine

同一批事件同时写入 Analytics Engine。

其主要用途是:

Event

Analytics Engine

Metrics / Reports

相比 R2 的原始数据存储,它更适合分析型查询。

因此系统形成两个数据出口:

                  Event

             ┌──────┴──────┐
             ▼             ▼
            R2       Analytics Engine
             │             │
        Historical      Analytics

二十一、第六层:数据维护层

R2 初始写入阶段采用“一条消息一个 JSONL 文件”。

例如:

raw/2026-09-16/10/

001.jsonl
002.jsonl
003.jsonl
...

这种方式适合高吞吐写入,但会产生大量小文件。

因此系统增加 Cron Compaction。


二十二、Compaction

每小时触发一次:

Cron


Select Completed Time Slot


Read Batch Files


Concatenate


compacted.jsonl


Delete Original Files

当前实现选择合并 2 小时前 的时间槽,从而避免当前仍在写入的数据与合并任务发生冲突。


二十三、为什么需要 Compaction

假设一天产生:

120 requests/hour

那么大量小对象会增加:

R2 List
R2 Get
Metadata
Query

等操作成本。

经过合并:

Before

00/
├── 001.jsonl
├── 002.jsonl
├── ...
└── 120.jsonl


After

00/
└── compacted.jsonl

原方案预计单日查询文件数量可以从约 2880 个降低到约 24 个。


二十四、第七层:查询层

Report 服务负责从 R2 查询历史事件。

查询流程:

Query


Resolve Date Range


Build File List


Concurrent R2.get()


Stream JSONL


Field Filter


Limit


Response

二十五、查询条件

查询支持:

date
date_start
date_end
ts_start
ts_end
event_name
任意字段
filter_mode
limit

其中:

filter_mode

支持:

AND
OR

字段匹配采用包含匹配,并且大小写不敏感。


二十六、并发读取

查询多个时间分区时,可以并发读取 R2:

              Query

        ┌───────┼───────┐
        ▼       ▼       ▼
       R2      R2      R2
        │       │       │
        └───────┼───────┘

             Stream

当前实现最多使用 50 个并发协程读取 R2。


二十七、Streaming Filter

查询不需要将整个文件一次性加载到内存。

采用:

R2 Object


Readable Stream


Parse Line


Filter

每读取一行就判断是否满足条件。

Line 1 → Match
Line 2 → Skip
Line 3 → Match
Line 4 → Skip

这样可以降低查询过程中的内存压力。


二十八、提前终止

当查询结果达到:

limit

后即可停止继续读取。

例如:

limit = 2000

查询:

Read

Match

Count = 2000

Stop

没有必要继续扫描剩余数据。

当前方案明确采用“命中 limit 后提前终止”的策略。


二十九、查询缓存

对于小规模查询,可以使用 Cloudflare Cache API。

当前策略:

limit <= 10,000


CF Cache

     TTL
    120s

大查询则不进入 Cache。

Small Query

Cache

Large Query

R2

缓存 Key 使用完整请求 URL。


三十、完整数据链路

将前面的各层组合起来:

                         Browser SDK


                       ┌────────────┐
                       │ Collect API│
                       └─────┬──────┘

                       Authentication

                         Decrypt

                         Validate

                       Geo Enrichment


                       ┌────────────┐
                       │    Queue   │
                       └─────┬──────┘


                       Queue Consumer

                    ┌────────┴────────┐
                    ▼                 ▼
                   R2          Analytics Engine


                Compaction


                Report API


                  Cache


                Admin / BI

三十一、代码模块与架构层映射

当前代码结构可以直接映射到架构层:

src/

├── index.ts
│   └── 接入层

├── auth.ts
│   └── 认证层

├── collect.ts
│   └── 数据采集层

├── consumer.ts
│   └── 异步处理层

├── scheduled.ts
│   └── 数据维护层

└── report.ts
    └── 查询层

这种映射非常直接,也与当前项目实际实现一致。


三十二、接口设计

目前主要提供两个接口:

MethodPath认证功能
POST/service_stats/v1/collectHMAC / API Key上报加密事件
GET/service_stats/v1/reportAPI Key查询事件数据

接口设计保持简单:

Collect
→ Write Path

Report
→ Read Path

写入和读取形成清晰的数据边界。


三十三、可靠性设计

整个系统的可靠性主要依赖三个机制。

Queue

负责:

Async
Retry
Decoupling

Consumer Idempotency

负责:

Duplicate Message

Idempotent Write

Compaction Window

负责:

Active Data

Wait

Stable Data

Compaction

三个机制共同保证数据链路稳定运行。


三十四、当前架构限制

大查询性能

当单次查询达到 100,000 条数据时:

R2 Read
 +
JSON Parse
 +
Field Filter

会产生明显 CPU 消耗。

当前方案已经识别到大查询可能受到 Cloudflare Worker CPU 限制,需要更高规格的运行模式。


Queue 延迟

Queue 是异步系统,因此存在秒级调度延迟。

因此:

Queue

R2

并不适合作为严格实时的数据查询来源。

如果需要更高实时性的分析,应优先使用 Analytics Engine。


Compaction 延迟

由于只处理已经稳定的数据,最近一段时间的数据可能仍然分散在多个批次文件中。

因此:

Recent Data
→ More Files

Historical Data
→ Compacted

查询近期数据的成本会高于历史数据。


三十五、测试设计

测试按照分层进行。

认证层

Expired Signature
Invalid Signature
Valid Signature

数据采集层

Invalid Payload
Modified Ciphertext
Invalid Event Count
Geo Enrichment

Queue 层

Duplicate Message
Retry
Idempotent Write

Compaction

Correct Time Window
Merge
Delete Source Files

Query

Date Filter
Cross-Date Query
Timestamp Filter
AND
OR
Limit

Cache

<= 10,000 → Cached
> 10,000  → Direct Query

原方案已经覆盖 HMAC、AES、Queue 幂等、Compaction、查询过滤、缓存以及并发读取等测试点。


三十六、总结

BI 埋点服务最终形成的是一条完整的数据处理流水线:

Collect

Authenticate

Decrypt

Validate

Enrich

Queue

Consume

Persist

Compact

Query

Cache

Report

其中每一层承担明确职责:

核心职责
接入层HTTP 路由与请求分发
认证层API Key / HMAC / Timestamp
数据采集层解密、校验、Geo 注入
异步处理层Queue 消费与幂等
数据存储层R2 / Analytics Engine
数据维护层定时 Compaction
查询层过滤、分页、缓存

最终实现:

采集与业务解耦、写入与查询解耦、实时分析与历史存储解耦。

这套架构的核心并不是单独使用某一个 Cloudflare 服务,而是利用 Workers + Queue + R2 + Analytics Engine 将同步请求、异步处理、历史存储和分析查询组合成一条完整的数据链路。

本页目录