- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本文以 Apache Beam 仓库中 .test-infra/metrics/sync/github/README.md 为核心,结合该目录下的同步脚本、GraphQL 查询、工具函数与 Docker 配置,完整讲解 Beam 社区如何将 GitHub 上的 Pull Request 与 Issue 数据持续同步进 PostgreSQL,为 metrics.beam.apache.org 上的 Grafana 社区指标大盘提供数据源。读完本文,你将掌握该同步器的容器化构建、本地运行、环境变量配置、数据表结构设计、增量同步机制与代码级实现细节,并可直接在本地复现整套运行流程。
一、背景:Beam 社区指标栈中的数据采集层
Apache Beam 的社区健康度(Community Metrics)需要持续观测两类外部数据源:Jenkins 上的 CI 构建结果,以及 GitHub 上的协作活动(PR、Issue、评审、提及等)。在 .test-infra/metrics/README.md 中明确指出,社区指标栈包含“从数据源(Jenkins 和 GitHub)摄取数据的 Python 脚本”和“Postgres 分析数据库”两部分;测试结果指标则另由 InfluxDB 时序数据库承载,最终两类指标统一呈现在 Grafana 大盘中。
本文聚焦其中的GitHub 数据摄取链路:位于 .test-infra/metrics/sync/github/ 目录下的syncgithub服务。它通过 GitHub GraphQL API(v4)拉取apache/beam仓库的 Pull Request 与 Issue 数据,清洗、结构化后写入 PostgreSQL,供 Grafana 查询展示。整个链路由 .test-infra/metrics/docker-compose.yml 中的syncgithub服务编排,而 .test-infra/metrics/sync/jenkins/README.md 对应的syncjenkins服务则承担 Jenkins 侧的数据同步,二者共同构成社区指标的数据采集层。
二、目录结构:GitHub 同步器的组成文件
.test-infra/metrics/sync/github/ ├── Dockerfile # 容器镜像定义(python:3.10-slim 基础镜像) ├── README.md # 本地运行与 lint 说明(本文核心文档) ├── sync.py # 主同步脚本:连接 DB、拉取 GitHub 数据、落库 ├── queries.py # GraphQL 查询定义(PR 与 Issue 两类查询) ├── ghutilities.py # GitHub 时间格式转换与 @提及提取工具 ├── sync_test.py # 针对 ghutilities 的单元测试 ├── requirements.txt # Python 依赖声明 └── github_runs_prefetcher/ # 独立的 GitHub Actions 工作流运行数据预取器其中github_runs_prefetcher是另一个独立子项目,负责把 GitHub Actions 工作流运行数据写入 CloudSQL,用于 Grafana 状态大盘与告警(详见其 README),与本文的 PR/Issue 同步器职责不同,注意区分。
三、本地运行:从构建镜像到执行同步
原文档给出了两个最核心的本地操作命令,此处结合源码逐条展开。
3.1 构建容器
首先需要在 .test-infra/metrics/sync/github/ 目录下构建 Docker 镜像(原文档步骤 1):
cd .test-infra/metrics/sync/github docker build -t syncgithub .Dockerfile 定义了镜像的构建方式:
FROM python:3.10-slim WORKDIR /usr/src/app COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt pylint yapf nose COPY . . CMD python ./sync.py关键点:
- 基础镜像为
python:3.10-slim,工作目录为/usr/src/app; - 除
requirements.txt中的运行时依赖外,还安装了pylint、yapf、nose三个开发/检查工具,说明该镜像既用于运行同步,也用于代码质量检查; - 镜像默认 CMD 直接执行
python ./sync.py,即不传参时容器启动即开始同步。
3.2 运行同步脚本
原文档的核心命令如下:
docker run -it --rm --name sync -v "$PWD":/usr/src/myapp -w /usr/src/myapp \ -e "DB_PORT=5432" \ -e "DB_DBNAME=beam_metrics" \ -e "DB_DBUSERNAME=admin" \ -e "DB_DBPWD=aaa" \ -e "GH_ACCESSTOKEN=<githubaccesstoken>" \ syncgithub python sync.py参数逐项说明:
| 参数 | 含义 | 对应源码读取位置 |
|---|---|---|
-it --rm --name sync | 交互式运行、退出即删除容器、命名容器为sync | — |
-v "$PWD":/usr/src/myapp | 把当前目录挂载进容器/usr/src/myapp,便于运行挂载目录内的sync.py | — |
-w /usr/src/myapp | 设定容器工作目录为挂载目录 | — |
-e DB_PORT=5432 | PostgreSQL 端口 | sync.py中DB_PORT = os.environ['DB_PORT'] |
-e DB_DBNAME=beam_metrics | PostgreSQL 数据库名 | DB_NAME = os.environ['DB_DBNAME'] |
-e DB_DBUSERNAME=admin | 数据库用户名 | DB_USER_NAME = os.environ['DB_DBUSERNAME'] |
-e DB_DBPWD=aaa | 数据库密码 | DB_PASSWORD = os.environ['DB_DBPWD'] |
-e GH_ACCESSTOKEN=<githubaccesstoken> | GitHub 个人访问令牌(PAT),用于 GraphQL 鉴权 | GH_ACCESS_TOKEN = os.environ['GH_ACCESS_TOKEN'] |
需要特别注意的是环境变量名的差异:文档中写的是GH_ACCESSTOKEN,而 sync.py 实际读取的是GH_ACCESS_TOKEN(下划线分隔)。本地直接运行python sync.py时,必须使用GH_ACCESS_TOKEN;若在 docker-compose 环境中则无需手动设置,由 compose 文件统一注入(见下文)。
另外,sync.py还会读取DB_HOST(第 43 行DB_HOST = os.environ['DB_HOST'])。原文档命令中没有设置它,这意味着本地运行场景下它必须预先存在于 shell 环境中(例如export DB_HOST=localhost),否则脚本会因KeyError启动失败。这一点是原文档未言明的隐含前提,本地复现时最容易踩坑。
3.3 运行 linter
原文档的第二条命令用于在容器内对sync.py执行 pylint 检查:
docker run -it --rm --name sync -v "$PWD":/usr/src/myapp -w /usr/src/myapp \ syncgithub pylint sync.py由于 Dockerfile 中已通过pip install ... pylint预装了 pylint,因此无需额外安装即可直接执行。同样的挂载与工作目录参数确保 pylint 能找到宿主机当前目录下的sync.py源文件。
四、主流程解析:sync.py 的同步逻辑
同步脚本 是整条链路的执行核心,其__main__入口(第 500 行起)展示了完整的运行循环:
print("Started.") initDbTablesIfNeeded() # 1. 检查并创建三张数据表 while True: # 2. 无限循环: if not probeGitHubIsUp(): # 探测 github.com:443 连通性 continue # 不通则跳过本轮 fetchNewData() # 通则执行增量拉取与落库 time.sleep(5 * 60) # 每 5 分钟一轮四个关键环节分述如下。
4.1 数据库连接与建表
initDBConnection()(第 94 行)使用psycopg2连接 PostgreSQL,连接串由DB_NAME、DB_USER_NAME、DB_HOST、DB_PORT、DB_PASSWORD五个环境变量拼装。连接失败时打印提示并sleep 60 秒后无限重试,这一设计让容器可以在数据库尚未就绪时安全启动(例如 docker-compose 中 Postgres 还在初始化)。
initDbTablesIfNeeded()(第 116 行)通过information_schema.tables检查三张表是否存在,不存在则按以下 DDL 建表:
PR 表gh_pull_requests:
create table gh_pull_requests ( pr_id integer NOT NULL PRIMARY KEY, author varchar NOT NULL, created_ts timestamp NOT NULL, first_non_author_activity_ts timestamp NULL, first_non_author_activity_author varchar NULL, closed_ts timestamp NULL, updated_ts timestamp NOT NULL, is_merged boolean NOT NULL, requested_reviewers varchar[] NOT NULL, beam_reviewers varchar[] NOT NULL, mentioned varchar[] NOT NULL, reviewed_by varchar[] NOT NULL )Issue 表gh_issues:
create table gh_issues ( issue_id integer NOT NULL PRIMARY KEY, author varchar NOT NULL, created_ts timestamp NOT NULL, updated_ts timestamp NOT NULL, closed_ts timestamp NULL, title varchar NOT NULL, assignees varchar[] NOT NULL, labels varchar[] NOT NULL )同步元数据表gh_sync_metadata(记录每次同步的游标时间):
create table gh_sync_metadata ( name varchar NOT NULL PRIMARY KEY, timestamp timestamp NOT NULL )注意 PR 表中的requested_reviewers、beam_reviewers、mentioned、reviewed_by以及 Issue 表中的assignees、labels均使用了 PostgreSQL 数组类型varchar[],用于容纳一人对多人(多个 GitHub 用户、多个标签)的度量维度。
4.2 增量同步游标:gh_sync_metadata
fetchLastSyncTimestamp(cursor, name)(第 164 行)按名称从元数据表读取上次同步时间戳,用作本轮查询的时间下界。对应地,updateLastSyncTimestamp(timestamp, name)(第 178 行)在每轮结束后以ON CONFLICT (name) DO UPDATE SET timestamp = excluded.timestamp的方式回写游标。
在 fetchNewData() 中,同步分为两个独立的游标:
- PR 同步游标名为
gh_pr_sync,Issue 同步游标名为gh_issue_sync; - 首次运行时元数据表为空,PR 走
fetchLastSyncTimestampFallback()(第 149 行),回退到datetime(year=1980, month=1, day=1),即从头全量拉取;Issue 则直接使用 1980-01-01 作为起点。源码中留有 TODO 注释,说明待gh_issue_sync行稳定存在后可移除回退逻辑。
4.3 GraphQL 查询与分页拉取
查询定义集中在 queries.py:
MAIN_PR_QUERY:通过 GitHubsearchAPI 查询apache/beam仓库的 PR,搜索条件为type:pr repo:apache/beam updated:><TemstampSubstitueLocation> sort:updated-asc,first: 100;MAIN_ISSUES_QUERY:结构类似,条件为type:issue repo:apache/beam updated:>... sort:updated-asc,first: 100。
两条查询都使用占位符<TemstampSubstitueLocation>,由 sync.py 的fetchGHData()在执行前替换为经过ghutilities.datetimeToGHTimeStr()格式化的时间字符串(格式%Y-%m-%dT%H:%M:%SZ,即 GitHub 标准时间格式)。
查询覆盖面很完整,PR 查询除了基础字段(number、author.login、createdAt、updatedAt、closedAt、merged、mergedAt、mergedBy.login、url、body)外,还嵌套拉取了:
- 最多 100 条评论(
comments,含作者、正文、创建时间); - 最多 50 条评审请求(
reviewRequests,含被请求评审人); - 最多 50 位 assignee(
assignees); - 最多 50 条评审记录(
reviews,含作者、正文、创建时间、状态)。
Issue 查询则拉取number、author.login、createdAt、closedAt、updatedAt、title、最多 50 位 assignee 与最多 10 个标签。
值得注意的细节:PR/Issue 的搜索均按updated:>过滤并sort:updated-asc(按更新时间升序),first: 100作为每页大小。但 sync.py 的fetchNewData()中的while resultsPresent循环并未真正消费pageInfo.endCursor进行翻页,而是依赖“更新游标推进”策略:每处理完一条记录,就把currTS更新为该记录的updatedAt(第 479–481 行),下一轮查询以更晚的时间为下界继续,从而在 API 单页 100 条的限制下通过多轮迭代完成全量追赶。这种“时间游标代替分页游标”的做法在该脚本中是刻意为之的简化实现。
4.4 数据提取与 upsert 落库
拿到 GraphQL 响应后,sync.py 通过一系列提取函数把节点数据转换为行值:
extractUserLogin()(第 210 行):用户节点可能缺失,缺失时返回"Unknown";extractRequestedReviewers()(第 216 行):从reviewRequests.edges提取被请求评审人登录名列表;extractMentions()(第 222 行):聚合 PR 正文、评论正文、评审正文中的@提及(经ghutilities.findMentions()用正则@(\w+)匹配,并过滤掉"username"字样);extractFirstNAActivity()(第 241 行):找出第一个由非作者发起的活动(评论、评审或合并)的时间戳与操作者,用于度量 PR 的“首次他人反馈等待时间”;若合并发生且合并者非作者,也会计入比较;extractBeamReviewers()(第 272 行):综合 assignees、reviewRequests、reviews 三类直接评审者信号,再加上正文/评论中通过特殊正则识别出的评审者标记,是逻辑最复杂的一环:beam_reviewer_regex = r'(@\w+).*?(?:PTAL|ptal|look)':匹配“@某人 ... PTAL/look”形式的评审请求;contrib_reviewer_regex = r'(?:^|\W)[Rr]\s*=:.+)':匹配R= @r1 @r2形式的贡献者评审标记;username_regex = r'(-?)(@\w+)':其中-前缀表示从评审者列表中移除该用户,实现“先加后减”的语义;- 最终通过
set去重并剔除作者本人(if r != author);
extractReviewers()(第 309 行):仅统计真正提交过 review 的作者(reviews.edges)。
随后extractRowValuesFromPr()/extractRowValuesFromIssue()把上述结果组装为行值数组,交给upsertIntoPRsTable()/upsertIntoIssuesTable()执行INSERT ... ON CONFLICT DO UPDATE的 upsert 写入(第 354、387 行),以pr_id/issue_id为冲突键,保证重复同步不产生脏数据。
五、环境变量全景:docker-compose 中的标准配置
除了本地docker run手动传参的方式,仓库还通过 docker-compose.yml 提供了标准化的服务编排。其中syncgithub服务的环境变量如下:
syncgithub: image: syncgithub container_name: beamsyncgithub build: context: ./sync/github dockerfile: Dockerfile environment: - DB_HOST=beampostgresql - DB_PORT=5432 - DB_DBNAME=beam_metrics - DB_DBUSERNAME=admin - DB_DBPWD=<PGPasswordHere> - GH_APP_ID=<GithubAppID> - GH_APP_INSTALLATION_ID=<GithubAppInstallationID> - GH_PEM_KEY=<GithubPemKey> - GH_NUMBER_OF_WORKFLOW_RUNS_TO_FETCH=30对比可发现:compose 环境额外注入了GH_APP_ID、GH_APP_INSTALLATION_ID、GH_PEM_KEY、GH_NUMBER_OF_WORKFLOW_RUNS_TO_FETCH等与GitHub App 鉴权及 Actions 工作流运行拉取相关的变量(这些由github_runs_prefetcher相关逻辑使用);而 sync.py 核心逻辑只需DB_*五件套与GH_ACCESS_TOKEN。若在 compose 全栈环境中运行,GH_ACCESS_TOKEN需要另行补充注入,否则sync.py会因读取不到该环境变量而退出。
完整的本地指标栈由postgresql(Postgres 9 系镜像、库名beam_metrics、用户admin)、influxdb(1.8.0,测试指标时序库)、grafana(含 PSQL/Influx 数据源配置与 JSON 数据源插件)与两个同步器组成;Grafana 大盘地址为http://localhost:3000,Postgres 映射到localhost:5432,InfluxDB 映射到http://localhost:8086。
六、可靠性设计与测试佐证
6.1 脚本的容错与自愈设计
从 sync.py 可以总结出若干工程化细节:
- DB 连接重试:连接失败时 sleep 60 秒重试,容忍数据库启动延迟;
- GitHub 连通性探测:
probeGitHubIsUp()(第 490 行)通过 TCP 连接github.com:443判断 GitHub 是否可达,不可达则跳过本轮,避免无谓的 API 调用与报错刷屏; - API 异常兜底:GraphQL 响应含
errors字段或 JSON 结构异常(如触发限流)时,打印错误并安全返回,等待下一轮 5 分钟后的重试; - 单条记录提取容错:某条 PR/Issue 提取失败时打印异常与
traceback后整体返回(第 467–472 行),防止脏数据半写入; - 幂等写入:所有落库均为 upsert,重复同步不会产生重复行;
- 游标持久化:同步进度存于
gh_sync_metadata表,容器重启后可从断点续传。
6.2 单元测试
sync_test.py 使用unittest与ddt数据驱动框架,覆盖了ghutilities.findMentions()的三种典型场景:
| 输入 | 期望输出 |
|---|---|
"sample text with mention @mention" | ["mention"] |
"Data without mention" | [] |
"sample text with several mentions @first, @second @third" | ["first", "second", "third"] |
该测试从侧面印证了findMentions()是同步链路中提取“提及”维度的基础能力,也是 sync.pyextractMentions()的底层依赖。注意测试中声明的test_findCommentReviewers目前只有占位实现,尚未完成,属于仓库中的已知未完成项。
七、数据流向总结与后续扩展
整条 GitHub 社区指标同步链路可归纳为:
GitHub GraphQL API (apache/beam 仓库) │ MAIN_PR_QUERY / MAIN_ISSUES_QUERY(按 updated 时间增量、每次 100 条) ▼ sync.py(ghutilities 提取提及 / 评审者 / 首次非作者活动等派生指标) │ INSERT ... ON CONFLICT DO UPDATE ▼ PostgreSQL(gh_pull_requests / gh_issues / gh_sync_metadata,库 beam_metrics) │ Grafana(kubeproxyuser_ro 只读用户查询) ▼ metrics.beam.apache.org 社区指标大盘对希望在本仓库基础上继续深入或二次开发的读者,建议按以下顺序阅读源码:
- .test-infra/metrics/README.md —— 了解整个 Beam 指标栈(社区指标 + 测试结果指标)的部署全貌;
- .test-infra/metrics/sync/github/sync.py —— 同步器主逻辑(连接、建表、增量拉取、提取、upsert、主循环);
- .test-infra/metrics/sync/github/queries.py —— 两条 GraphQL 查询的完整字段定义,如需新增维度(例如评论数、文件改动数)在此扩展;
- .test-infra/metrics/sync/github/ghutilities.py —— 时间格式互转与 @提及提取;
- .test-infra/metrics/docker-compose.yml —— 标准化的本地部署与参数注入。
需要说明的适用前提:本文所有命令、表结构与环境变量均以当前仓库快照为准;sync.py依赖DB_HOST与GH_ACCESS_TOKEN两个文档未完整写明的环境变量,本地直接运行时务必先行注入;GitHub GraphQL API 存在速率限制,同步器通过 5 分钟轮询与时间游标设计来规避,若需更高拉取频率应在充分评估配额的前提下修改time.sleep(5 * 60)的间隔。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam社区贡献指南:从Issue到PR的完整流程
Apache Beam社区贡献指南:从Issue到PR的完整流程 作为Apache Beam(Apache软件基金会旗下的统一批处理和流处理编程模型)的贡献者,
批处理流处理大数据缺陷报告
缺陷报告 环境信息 SkyWalking版本: 9.7.0 部署方式: 容器化/Docker Compose 操作系统: Linux Ubuntu 22.04
可观测性后端微服务云原生GraphQL Playground社区贡献指南:从Issue到PR
GraphQL Playground社区贡献指南:从Issue到PR 作为开源项目,GraphQL Playground的发展离不开社区贡献。本文将详细介绍从发
开发工具后端API设计
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考