Hatchet gRPC 压缩测试套件:跨 Go/TypeScript/Python SDK 的网络流量基准与源码解析
【免费下载链接】hatchet🪓 An orchestration engine for background tasks, AI agents, and durable workflows项目地址: https://gitcode.com/GitHub_Trending/ha/hatchet
本文围绕 Hatchet 仓库中的 gRPC 压缩测试套件(hack/dev/compression-test/)展开:它如何以 Docker 容器化方式驱动 Go、TypeScript、Python 三个 SDK 客户端,通过docker stats精确测量每个 SDK 的进出流量,并对比“压缩开启/关闭”两种状态下的带宽差异。读完后,你将掌握该测试套件的完整构建与运行流程、所有环境变量与测试参数的取值含义,以及压缩开关在 Go SDK gRPC 客户端中的实际生效位置,从而可以复制这一方法论来量化自己项目中传输层的任何改动。
1. 套件定位:测量“压缩前后”的网络流量差
套件说明文档将其定义为“test gRPC compression across Go, TypeScript, and Python SDKs”的脚本与配置集合。核心思路是控制变量对比:
- Baseline(基线):
main分支构建的镜像,代表无压缩状态; - With Compression(压缩):当前分支构建的镜像,代表启用压缩状态。
两类镜像分别打标签{sdk}-{state}-compression,具体为:
go-disabled-compression/go-enabled-compressiontypescript-disabled-compression/typescript-enabled-compressionpython-disabled-compression/python-enabled-compressionengine-disabled-compression/engine-enabled-compression
其中 engine 镜像是可选的——README 明确说明引擎可以自行单独管理(# docker build -t engine-... -f Dockerfile .仅作注释形式的提示)。
标准化的测试参数在文档中固定为:
| 参数 | 取值 |
|---|---|
| 持续时间 Duration | 60 秒 |
| 事件速率 | 10 events/s |
| Payload 大小 | 100KB |
| DAG 步骤数 | 1 |
| 事件扇出 eventFanout | 1 |
按 QUICKSTART.md 的描述,每个 SDK 会在 60 秒内发出 600 个事件(10 事件/秒),每个事件携带 100KB payload,最终由报告脚本计算压缩带来的带宽节省比例。
2. 压缩在 Go SDK 中的实际生效位置
测试对比的“compression”并不是测试脚本本身实现的,而是被测 SDK 代码内的 gRPC 传输能力。以仓库内 Go 客户端为例,pkg/client/client.go 中的拨号逻辑显示:
grpcOpts := []grpc.DialOption{ grpc.WithTransportCredentials(transportCredentials), grpc.WithKeepaliveParams(keepAliveParams), } if !opts.disableGzipCompression { grpcOpts = append(grpcOpts, grpc.WithDefaultCallOptions(grpc.UseCompressor("gzip"))) opts.l.Info().Msg("gzip compression enabled for gRPC client") } else { opts.l.Info().Msg("gzip compression disabled for gRPC client") }即当disableGzipCompression为 false 时,客户端会为默认调用选项挂上grpc.UseCompressor("gzip"),并在日志中打印启用状态;gzip codec 则通过包顶层的空白导入_ "google.golang.org/grpc/encoding/gzip"完成注册(pkg/client/client.go)。同样的模式也出现在 pkg/client/dispatcher.go 与 pkg/client/v1/grpc-client.go。
这解释了测试套件的分支对比逻辑:main分支与压缩分支构建出的镜像差异,正体现在这类开关上。测试脚本本身只做流量测量,不干预压缩行为。
3. 构建镜像:三个 Dockerfile 与一键脚本
3.1 手动构建(README 原始流程)
README 强调:所有docker build命令必须从仓库根目录执行,而不是在测试目录下执行,因为构建上下文需要覆盖整个仓库。
Step 1:切到main分支,构建基线镜像:
git checkout main cd /path/to/hatchet # Build Go SDK docker build -t go-disabled-compression -f hack/dev/compression-test/Dockerfile.client-go . # Build TypeScript SDK docker build -t typescript-disabled-compression -f hack/dev/compression-test/Dockerfile.client-ts . # Build Python SDK docker build -t python-disabled-compression -f hack/dev/compression-test/Dockerfile.client-python . # Build Engine (optional - you may manage this separately) # docker build -t engine-disabled-compression -f Dockerfile .Step 2:切到压缩分支,用同样的三个 Dockerfile 构建*-enabled-compression镜像。
3.2 一键构建脚本
scripts/build_all.sh 把上述流程脚本化:它先自行定位仓库根目录(REPO_ROOT="$(cd "$TEST_DIR/../../.." && pwd)"),再依次构建三个镜像,失败即退出。用法(见 QUICKSTART.md):
cd /path/to/hatchet ./hack/dev/compression-test/scripts/build_all.sh enabled参数enabled/disabled决定镜像标签后缀,默认enabled。
3.3 三个 Dockerfile 分别构建什么
- Dockerfile.client-go:两阶段构建。builder 阶段基于
golang:1.26-alpine,先COPY go.mod go.sum触发go mod download缓存层,再只拷贝cmd/hatchet-loadtest/、pkg/、internal/、api/、api-contracts/,以CGO_ENABLED=0 GOOS=linux go build -tags=load编译出hatchet-loadtest二进制;运行阶段仅alpine+ca-certificates,默认CMD即为一次完整的 loadtest 调用(--events 10 --duration 60s --payloadSize 100kb --dagSteps 1 --eventFanout 1 --slots 100)。 - Dockerfile.client-ts:builder 阶段基于
node:20-alpine,用 pnpm 安装sdks/typescript依赖并执行pnpm run tsc:build编译 SDK,同时把测试文件 tests/typescript_test.ts 拷入镜像;运行阶段全局安装ts-node typescript,设置NODE_PATH=/app/dist,以ts-node -r tsconfig-paths/register直接运行测试脚本。 - Dockerfile.client-python:builder 阶段基于
python:3.11-slim,用poetry build从sdks/python/打出 wheel;运行阶段pip install *.whl安装本地构建的 SDK,并拷贝 tests/python_test.py 作为入口。
三者共同点:被测对象都是当前分支源码现场构建的 SDK,保证“enabled/disabled”两组镜像之间的差异只来自压缩相关代码。
4. 运行前准备:环境变量与前置条件
4.1 前置条件
- 引擎必须先独立运行:这些脚本不管理引擎生命周期,需自行启动 Hatchet engine,并保证它可通过
HATCHET_CLIENT_HOST_PORT指定的地址访问; - 设置必需环境变量:
export HATCHET_CLIENT_TOKEN="your-token-here" export HATCHET_CLIENT_HOST_PORT="localhost:7070" # gRPC address where your engine is running- 可选环境变量(含 README 标注的默认值):
export HATCHET_CLIENT_SERVER_URL="http://localhost:8080" # HTTP server URL (for Go SDK, defaults to http://localhost:8080) export HATCHET_CLIENT_API_URL="http://localhost:8080" # API URL (for TypeScript SDK, defaults to http://localhost:8080) export HATCHET_CLIENT_NAMESPACE="compression-test" # Namespace (optional, defaults to compression-test)两个补充说明(同样来自 README):
HATCHET_CLIENT_TENANT_ID会自动从 token 中提取,无需手动设置;- 脚本使用
host网络模式访问方式,容器可直连宿主机上的引擎。
4.2 源码中的变量消费链路
docker-compose.yml 展示了这些变量如何流入三个客户端容器(client-go/client-typescript/client-python)。每个服务都通过extra_hosts: "host.docker.internal:host-gateway"解析宿主机地址,且HATCHET_CLIENT_TLS_STRATEGY=none固定关闭 TLS,确保流量差异只来自压缩本身。
scripts/run_test.sh 则体现了两层防御:
HATCHET_CLIENT_TOKEN缺失时直接报错退出;HATCHET_CLIENT_HOST_PORT缺失时,自动通过docker run --rm alpine getent hosts host.docker.internal探测 Docker 网关 IP(回退值192.168.65.254),默认拼成${GATEWAY_IP}:7070——这是针对 macOS Docker Desktop 规避 IPv6 解析问题的处理。
运行前还会执行 scripts/setup.sh:创建results/baseline与results/enabled结果目录,并检查docker与docker compose是否可用。
5. 执行测试:全量与单 SDK 两种入口
5.1 Quick Start(README 原始流程)
cd hack/dev/compression-test # Run setup (creates network and directories) ./scripts/setup.sh # Run all baseline tests ./scripts/run_all_tests.sh disabled # Switch engine to compression version, then run compression tests ./scripts/run_all_tests.sh enabled # Generate comparison report ./scripts/generate_report.sh注意run_all_tests.sh disabled/enabled切换的是客户端镜像({sdk}-{state}-compression);README 提示在跑enabled组前需自行把引擎切到压缩版本。
5.2 单 SDK 测试
# Test Go SDK (baseline / compression) ./scripts/run_test.sh go disabled ./scripts/run_test.sh go enabled # Test TypeScript SDK ./scripts/run_test.sh typescript disabled ./scripts/run_test.sh typescript enabled # Test Python SDK ./scripts/run_test.sh python disabled ./scripts/run_test.sh python enabledscripts/run_all_tests.sh 的批量执行顺序是python → typescript → go,每次之间 sleep 5 秒;EVENTS_COUNT为可选第二参数,默认 10。
5.3 单测脚本内部做了什么
run_test.sh的核心流程(源码逐段印证):
- 校验镜像
${SDK}-${STATE}-compression是否存在(docker image inspect),不存在则提示先构建; - 清理同名旧容器
hatchet-client-${SDK},导出环境变量后以docker-compose run -d启动对应服务; - 速率换算:脚本固定
EVENTS_PER_SECOND=10,由EVENTS_COUNT反推TEST_DURATION(向上取整、至少 1 秒),并设置TEST_WAIT=时长+5s。这里有一个值得注意的语义细节(源码注释原文):Go 端--events是每秒速率而非总量,所以要duration = count / rate; - 后台并行启动 scripts/monitor_network.sh 与实时日志流(
docker logs -f --tail 10); - 轮询等待容器退出,超时阈值:TypeScript 180 秒,其余 120 秒;
- 容器结束后读取
${SDK}_network.log.summary,打印 RX/TX/Total 字节数。
docker-compose.yml中各 SDK 的启动命令也值得对照:
- Go:
./hatchet-loadtest loadtest --events ${TEST_EVENTS_RATE} --duration ${TEST_DURATION} --payloadSize 100kb --dagSteps 1 --eventFanout 1 --slots 100 --wait ${TEST_WAIT}; - TypeScript:
ts-node -r tsconfig-paths/register typescript_test.ts,事件数由TEST_EVENTS_COUNT控制; - Python:
python python_test.py,同样读取TEST_EVENTS_COUNT。
其中 Go 侧的负载生成器是仓库内的 cmd/hatchet-loadtest,其 flag 定义明确了各参数语义:--events(events per second,默认 10)、--duration(总时长)、--wait(等待事件完成的时间)、--payloadSize(payload 大小,默认 0kb)、--dagSteps(DAG 步数)、--eventFanout(事件扇出)、--slots(worker 槽位数)。
5.4 TS/Python 测试脚本的工作负载
两个脚本构建完全同构的工作负载,保证跨语言可比:
- 100KB payload:python_test.py 与 typescript_test.ts 都是“100 个 key,每个值为 1000 字符的 chunk”(
chunk = "a" * 1000,循环 100 次),整体约为 100KB 的高重复性 JSON 字典——这类数据对 gzip 压缩最友好,正好用来放大压缩收益; - 事件推送:以 100ms 间隔向事件名
compression-test:event推送{id, createdAt, payload}; - 单步工作流:workflow 名含压缩状态(如
enabled-python),on_events=["compression-test:event"],任务step1只做读取与回显; - Worker 规格:
slots=100,与 Go 端--slots 100对齐; - 自终止:Python 端用
threading.Timer在duration + 15s后向自身发SIGTERM关闭 worker;TypeScript 端在事件发完后再等待 10 秒处理缓冲,然后带 10 秒超时竞速调用worker.stop(),最后process.exit(0)。
6. 网络流量测量:基于 docker stats 的 NetIO 差分
scripts/monitor_network.sh 是量化核心,其机制:
- 初始采样:每 5 秒执行一次
docker stats --no-stream --format "{{.NetIO}}" <container>,取容器累计的 NetIO 计数;容器刚启动时docker stats可能返回0B / 0B或-- / --,脚本会重试最多 10 次; - 单位解析:
parse_stats_to_bytes()用正则匹配1.2MB / 3.4MB这类输出(兼容 B/KB/MB/GB/TB 与 KiB/MiB/GiB/TiB 两套单位,并统一大小写),再用awk换算成纯字节; - 差分计算:测试结束时的累计 NetIO 减去初始采样,得到本次测试窗口内的
TOTAL_RX/TOTAL_TX(注意 docker stats 的 NetIO 是累计值而非速率,因此必须做差分); - 结果落盘:追加采样序列到
${OUTPUT_FILE},最终把三行摘要写入${OUTPUT_FILE}.summary:
RX_BYTES=$TOTAL_RX TX_BYTES=$TOTAL_TX TOTAL_BYTES=$TOTAL_BYTES脚本还注册了trap 'handle_exit' TERM INT:即使被run_test.sh中途终止,也会尝试读取最后一次有效统计并补写 summary,避免结果丢失。bc缺失时自动降级为awk计算。
7. 结果与报告:压缩收益的百分比对比
结果文件位于results/目录:
results/baseline/(setup.sh 实际创建为results/disabled一类状态目录)——基线测试结果,每个 SDK 一份${SDK}_network.log+${SDK}_network.log.summary;results/compressed/(对应results/enabled)——压缩测试结果。
scripts/generate_report.sh 先调用collect_results.sh聚合出results/aggregated_results.txt,再读取GO_BASELINE/GO_COMPRESSED、TYPESCRIPT_*、PYTHON_*、BASELINE_TOTAL/COMPRESSED_TOTAL等变量,用bc计算每个 SDK 与总量的下降百分比,生成results/compression_report.txt。报告结构为:
======================================== Compression Test Results ======================================== Baseline (No Compression): Go SDK / TypeScript SDK / Python SDK / Total With Compression: 各 SDK 压缩后流量 + (x% reduction) Total + (x% reduction) Bandwidth Savings: Total Saved: ... Reduction: x% ======================================== Detailed Breakdown ======================================== (每个 SDK 的 Baseline / Compressed / Savings 三段明细)由于 RX(拉取任务、流式回调)与 TX(事件推送、注册)都被计入TOTAL_BYTES,该报告反映的是客户端与引擎之间全部 gRPC 流量的净变化,而非单一方向。
8. 复现要点小结
- 镜像构建必须发生在仓库根目录,且 disabled/enabled 两组镜像必须分别来自
main分支与压缩分支的现场构建; - 引擎独立部署,客户端容器经
host.docker.internal/host 网络直连,TLS 关闭(HATCHET_CLIENT_TLS_STRATEGY=none); - 负载三要素固定:10 events/s、100KB 高重复 payload、单步 DAG,worker 一律 100 slots;
- 测量只依赖
docker stats的累计 NetIO 差分,与容器内 SDK 实现无关,因此对三种语言 SDK 的口径一致; - 报告输出的 reduction 百分比基于
bc精确计算,results/compression_report.txt即为最终交付物。
整套脚本的分工可以归纳为:setup.sh建目录 →build_all.sh建镜像 →run_all_tests.sh/run_test.sh驱动容器 →monitor_network.sh采流量 →collect_results.sh+generate_report.sh出报告。所有环节均可在 hack/dev/compression-test/ 目录内逐一查看源码,配合上文引用的 SDK 与 loadtest 源码路径,即可完整复现并改造这套压缩(或任何传输层改动)的 A/B 流量基准。
【免费下载链接】hatchet🪓 An orchestration engine for background tasks, AI agents, and durable workflows项目地址: https://gitcode.com/GitHub_Trending/ha/hatchet
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考