ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

PostHog Temporal 编排实战指南:从 Workflow/Activity 概念到超时、心跳、重试与测试模式

2026/9/13 7:50:36 拓冰建站 浏览量
PostHog Temporal 编排实战指南:从 Workflow/Activity 概念到超时、心跳、重试与测试模式 PostHog Temporal 编排实战指南从 Workflow/Activity 概念到超时、心跳、重试与测试模式【免费下载链接】posthog:hedgehog: PostHog is the leading platform for building self-driving products. Our developer tools – AI observability, analytics, session replay, flags, experiments, error tracking, logs, and more – capture all the context agents need to diagnose problems, uncover opportunities, and ship fixes. Steer it all from Slack, web, desktop, or the MCP.项目地址: https://gitcode.com/GitHub_Trending/po/posthog本文基于 PostHog 仓库中的posthog/temporal/README.md及配套源码编写系统讲解 PostHog 如何使用 Temporal 构建批量导出、数据仓库同步、AI 对话处理等产品的调度与编排层。读完后你将掌握 Temporal 的 Workflow/Activity/Schedule 核心概念、活动超时与心跳的设计原则、Heartbeater/ShutdownMonitor等 PostHog 自研抽象的用法以及如何将 workflow 注册到 worker、触发部署并在本地开发与测试。为什么 PostHog 选择 TemporalPostHog 使用 Temporal 驱动多个产品和特性全部的 Batch exports批量导出数据仓库data warehouse的同步AI agent 对话的处理数据仓库模型物化materialization以及更多……这些产品都需要一个具备幂等执行、细粒度错误管理、可观察执行历史的调度与编排系统并且能横向扩展到大量 worker。Temporal 的价值在于它以看起来像本地代码的抽象来处理分布式执行让你可以专注于编写应用逻辑。但代价同样真实——Temporal 会带来维护负担与额外复杂度原因一部分在于分布式系统本身很难另一部分在于 Temporal 的抽象并非完全无泄漏。因此官方建议在采用之前先评估你应用对可靠性、幂等性和可观测性的需求与维护成本和复杂度做权衡。一旦决定使用 TemporalPostHog 在posthog/temporal/下沉淀了大量公共工具来降低这套复杂度新特性可以直接复用。提示DuckLake 复制copyworkflow 的配置说明在 products/managed_warehouse/backend/README.md 中。Temporal 核心概念Workflow像 Airflow DAG 一样的代码化编排Workflow 是一串通常存在依赖关系的 activity 序列。如果你熟悉 Apache Airflow可以把 Temporal workflow 类比为 Airflow DAG。Workflow 是代码PostHog 使用 Temporal Python SDKtemporalio以 Python 类来编写 workflowSDK 也支持其他语言。具体而言workflow 是一个用temporalio.workflow.defn装饰的类其中有一个方法被temporalio.workflow.run装饰——该方法即 workflow 的入口。方法内部实现业务逻辑通过workflow.execute_activity调用 activity并聚合结果后返回。Workflow 代码运行在 Temporal 自实现的asyncio事件循环中因此必须是确定性的。最常见的限制是任何 I/O 调用都不能在 workflow 中直接执行必须委托给 activity。一个最小示例import asyncio from temporalio import activity, workflow from temporalio.client import Client activity.defn async def hello_world_activity() - str: return Hello world! workflow.defn class HelloWorldWorkflow: workflow.run async def run(self) - str: return await workflow.execute_activity(hello_world_activity) async def main(): client await Client.connect(localhost:7233) result await client.execute_workflow(HelloWorldWorkflow.run, idmy-workflow-id, task_queuemy-task-queue) print(fResult: {result}) if __name__ __main__: asyncio.run(main())提示在整个仓库中搜索workflow.defn装饰器可以找到大量真实的 workflow 实现。Activity三种并发形态的单次操作Activity 表示 workflow 执行的单次操作类比 Airflow 中的 task。与 workflow 一样activity 也是代码Python SDK 中用temporalio.activity.defn装饰函数即可。Activity 有三种类型协程async def函数同步多线程函数在线程池中运行同步多进程函数在进程池中运行。使用 Python SDK 时协程是推荐的默认选择——性能更好且无需额外配置。但如果某个产品需要执行不适合asyncio的工作例如 CPU 密集型计算activity 也可以写成同步函数根据 worker 配置运行在线程池或进程池中。如果 activity 容易遇到临时性故障可以配置 Temporal 自动重试。Temporal service编排器与持久化Temporal service 是所有 workflow 和 activity 的编排器。它保证 workflow 的持久性durability——即使 worker 崩溃workflow 依然存续。原理是 service 会记录 workflow 的执行进度这也使得 activity 失败后的自动重试成为可能。开发者需要时可以检查这份执行历史来调试 workflow。PostHog 依赖 Temporal Cloud 托管 Temporal service本地开发栈则自带一个本地 service。Worker按 task queue 消费工作Workflow 和 activity 运行在 Temporal worker 中。Temporal service 根据代码中的配置把 workflow 和 activity 的执行分配到不同的 task queueworker 轮询这些队列来执行代码。Worker 同样由 Temporal Python SDK 提供并在代码中配置。注意一个 worker 只能监听单个队列。PostHog 的 worker 配置大多集中在posthog/temporal/包中尤其是 worker.py 模块生产与本地的启动入口是 start_temporal_worker.py。Schedule定时启动 workflow 的四要素Temporal schedule 由四部分定义action、spec、state、policy。Actionschedule 能执行的动作永远是启动一个 workflow。使用temporalio.client.ScheduleActionStartWorkflow配置运行哪个 workflow、参数、ID、task queue 和 workflow 的重试策略。注意配置的 workflow ID 更像是一个ID 前缀——Temporal schedule 会在 ID 末尾追加 workflow 的开始时间。Spectemporalio.client.ScheduleSpec定义 schedule 的运行间隔。start_at/end_at决定 schedule 本身的起止——start_at之前不产生运行end_at之后不产生运行两者都可为None即永久运行直到被暂停或删除。警告Temporal 会在end_at过后的一段时间自动清理schedule。通常无需设置end_at但如果设置了请意识到这意味着 schedule 会被 Temporal service 永久删除且无法恢复。间隔有两种定义方式intervals参数接受datetime.timedelta集合决定每次运行之间的增量或calendars参数按日期成分如day_of_month、day_of_week匹配。重要基于日历的表达式默认按 UTC 解释可用time_zone_name参数覆盖。Spec 还可以包含jitter为每次运行的开始时间加上随机值。这对多个使用相同 spec 的 schedule 非常有用——避免所有运行同时开始从而防止对相关服务形成惊群thundering herd效应。Statetemporalio.client.ScheduleState仅表示 schedule 是否被暂停。被暂停的 schedule 在恢复前不会执行任何动作。Policytemporalio.client.SchedulePolicy决定是否允许运行 ID 重叠对 backfill 场景有用以及失败时是否暂停 schedule。完整示例import uuid from temporalio.common import RetryPolicy from temporalio.client import ( Schedule, ScheduleActionStartWorkflow, ScheduleIntervalSpec, ScheduleOverlapPolicy, SchedulePolicy, ScheduleSpec, ScheduleState, ) from posthog.temporal.common.schedule import a_create_schedule from products.my_product.workflows import my_workflow, MyWorkflowInputs my_inputs MyWorkflowInputs(...) my_retry_policy RetryPolicy(...) my_task_queue general-purpose-task-queue workflow_id str(uuid.uuid4()) schedule Schedule( actionScheduleActionStartWorkflow( workflowmy_workflow, argsmy_inputs, idworkflow_id, task_queuemy_task_queue, retry_policymy_retry_policy, ), specScheduleSpec( intervals[ScheduleIntervalSpec(every...)], jitter..., ), stateScheduleState(), policySchedulePolicy(overlapScheduleOverlapPolicy.ALLOW_ALL), )PostHog 在 common/schedule.py 中把这层封装成了同步/异步两套工具函数create_schedule/a_create_schedule创建、update_schedule/a_update_schedule更新支持keep_tz保留时区、pause_schedule/unpause_schedule、trigger_schedule、delete_schedule、schedule_exists通过 describe 判断捕获RPCStatusCode.NOT_FOUND以及trigger_schedule_buffer_one以BUFFER_ONE重叠策略触发。代码库中到处能看到这些工具的使用例如 batch exports 每创建一个批量导出就新建一个 schedule其他产品则用 Django management command 统一初始化所有 schedule——具体模式取决于产品自身结构。用 Temporal 实现一个产品或特性第一步是判断 Temporal 是否合适如果你的应用/特性需要把工作卸载给专门的 worker、要求持久性和幂等性、需要错误处理与重试管理工具并且愿意接受复杂度与运维成本的上升那么 Temporal 可能是合适的工具。选定之后posthog/temporal/common/下的公共抽象会帮你把新增的复杂度控制在小范围内。决策一选择 activity 的并发模型Temporal worker 同时运行多个 workflow 和 activity因此必须为 activity 选择并发模型asyncio、多线程、多进程。官方立场提示一个 workflow 可以为不同的 activity 使用不同的并发模型先执行协程 activity再执行多进程 activity。但为了保持简单推荐为你的所有 activity 找到单一可行的并发模型。提示asyncio 通常是最佳选择。Asyncio不要阻塞事件循环编写 asyncio 代码最重要的规则DO NOT BLOCK事件循环。Asyncio 最适合 I/O 密集的代码网络、数据库请求但前提是这些请求使用非阻塞原语完成。否则 worker 里其他任务无法并发执行asyncio 的性能优势尽失而且同一事件循环还承载 worker 中的其他 activity阻塞会波及它们最终导致超时。Asyncio 在 PostHog 单体的其他部分尚未大规模普及这意味着很多现成代码不能直接搬进 activity。特别是Django model 在 PostHog 其他位置使用的常规方法调用是阻塞的。因此从单体其他部分引入代码时通常需要一些改造工作优先级如下目标库已支持 asyncio 且有可直接替换的方法。例如 Django model 提供在方法名前加a的异步版本MyModel.objects.get(...)变为await MyModel.objects.aget(...)——但并非整个 model API 都支持需查当前 Django 版本文档。库本身不支持 asyncio但存在替代品。例如requests是阻塞的可换aiohttp、httpxAPI 非常接近boto3有aioboto3Kafka 有提供非阻塞生产/消费类的aiokafka。都不行时把阻塞代码丢进线程池concurrent.futures.ThreadPoolExecutor或直接asyncio.to_thread——Python 在 I/O 操作时会释放 GIL把代码发到别的线程即可避免阻塞事件循环。同理CPU 密集型阻塞代码可尝试concurrent.futures.ProcessPoolExecutor。实在无解就用 asyncio 库和原语重新实现。代码改造完成后它就能在 Temporal worker 中与其他协程并发协作。提示改造为 asyncio 后就能应用真正的asyncio 模式而不只是到处加await一串顺序请求可以用 task asyncio.gather或asyncio.TaskGroup并发化进度更新可以作为后台 task 让主流程继续数据到达即可用asyncio.Queue以生产者-消费者模式流式处理。多线程不推荐但可行不推荐同步多线程 activity——任何线程需求都可以在异步 activity 中用asyncio.to_thread满足线程池。如果你有充分理由坚持多线程只需把 activity 函数写成同步函数因为PostHog 的 worker 默认就配置了线程池执行器。但要注意创建线程比 asyncio 消耗更多资源且 PostHog 基于 asyncio 构建的许多上层抽象如Heartbeater需要额外的同步包装才能在同步代码中使用——例如用 heartbeat_sync.py 中的HeartbeaterSync替代Heartbeater。GIL 的考量同样存在除非线程代码包含能释放 GIL 的 I/O 操作否则其他 activity 无法运行这个并发模型就失去了意义。决策二为 activity 设置超时Temporal 允许对 activity 施加多种超时Schedule-to-close从 activity 任务进入队列那一刻起计时Start-to-close从 worker 开始执行 activity 那一刻起计时Schedule-to-start从 activity 被 worker 领取所花的时间计时Heartbeat基于心跳之间的间隔计时。每个 activity必须至少定义 schedule-to-close 和 start-to-close 二者之一README 推荐后者。这就是 Temporal service 在 worker 崩溃后恢复的方式超时到期后service 重新下发该 activity让存活的 worker 接手。决策三为长任务 activity 发送心跳对长任务而言仅靠上述两种超时不够。设想一个预期 1 小时完成的 activitystart_to_close设为 1 小时。如果 worker 在 activity 刚开始时崩溃你要等将近整整 1 小时超时到期service 才会重新调度——这 1 小时完全被浪费而下一小时的 workflow 又可能已经启动造成积压。这就是为什么任何长任务 activity 都强烈建议用心跳 心跳超时activity 通过心跳告知 service 自己还活着一旦停止心跳超过心跳超时时长service 就判定 worker 大概率已崩溃可以立即重试而无需等满 start-to-close。使用posthog.temporal.common.heartbeat.Heartbeater类实现心跳非常容易作为上下文管理器包裹你的长任务即可。从源码看Heartbeater 的默认factor120即它会在heartbeat_timeout / 120的间隔上定期发出心跳留足余量避免误判。除了常规心跳任务它还会启动第二个任务检测到 worker 关闭activity.wait_for_worker_shutdown()时立即把最新的 heartbeat details 刷出去——这对跨 worker 恢复进度至关重要。退出上下文时会取消这两个任务并补发最后一次心跳。此外还有一个带存活度追踪的LivenessHeartbeater子类以及可序列化的HeartbeatDetails数据类体系如DataImportHeartbeatDetails携带endpoint和分页cursor配合from_activity类方法可在重试后从activity.info().heartbeat_details恢复进度。心跳可以携带任意附加信息heartbeat details会被持久化到 Temporal 并可稍后取回从而支持用 heartbeat 做进度追踪但要注意心跳投递是缓冲的、不保证送达——一般追踪目的没问题精确控制请另寻他处。决策四配置重试策略Activity 失败后可能被重试行为由RetryPolicy决定最大尝试次数、重试间隔可实现指数退避等常见策略以及一个不可重试错误列表——列表中的错误即使还有重试次数也不会再试适合定义重试也解决不了的致命错误。重要Temporal 按异常类名匹配错误大概出于序列化/反序列化的考虑——一切皆分布式。设置non_retryable_error_types时必须列出你想排除的每一个异常的具体类名而不能写类层次中的公共祖先。默认情况下 activity 会永远重试。警告永远就是字面意义的永远。请认真考虑能否接受 activity 最终失败如果不能就配置告警及时发现并人工处理那些永远不会成功的 activity。提示测试中务必设置max_retries。下面这个综合示例覆盖了心跳、超时与重试的推荐写法仓库中搜索可找到更多import asyncio import datetime as dt from temporalio import activity, workflow from temporalio.common import RetryPolicy from posthog.temporal.common.heartbeat import Heartbeater class FatalError(Exception): ... class UserFatalError(FatalError): ... class InternalFatalError(FatalError): ... activity.defn async def short_activity() - None: # 短超时场景不需要心跳。 await do_short_work() activity.defn async def long_running_activity() - None: # 长任务 activity 必须心跳。 async with Heartbeater(): await do_long_work() workflow.defn class HelloWorldWorkflow: worklfow.run async def run(self) - None: await workflow.execute_activity( short_activity, start_to_close_timeoutdt.timedelta(seconds60), # 不设 maximum_attempts 意味着永远重试。 retry_policyRetryPolicy( initial_intervaldt.timedelta(seconds2), maximum_intervaldt.timedelta(seconds10), non_retryable_error_types[InternalFatalError], ), ) await workflow.execute_activity( long_running_activity, start_to_close_timeoutdt.timedelta(hours2), heartbeat_timeoutdt.timedelta(seconds30), retry_policyRetryPolicy( initial_intervaldt.timedelta(seconds2), backoff_coefficient2.0, maximum_intervaldt.timedelta(seconds64), maximum_attempts10, # 列出所有可能的错误类型而不是公共祖先。 non_retryable_error_types[InternalFatalError, UserFatalError], ), )决策五日志Write Produce 双通道与 PostHog 其余部分一致Temporal 的日志使用 structlog 中的Logger/ProduceOnlyLogger/WriteOnlyLogger与get_logger工厂函数Write日志写到 stdout供内部日志解析器与监控系统采集Produce日志投递到 Kafka随后被 ClickHouse 消费进log_entries表使日志可以直接在 PostHog 里查询——例如直接向用户展示。典型应用是 batch exports 的日志 tab把调试信息呈现给用户帮助他们修复配置错误。默认从structlog.get_logger拿到的 logger 会同时做 Write 和 Produce。注意Produce 模式需要额外配置才能匹配log_entries表结构上下文的某处必须设置team_id必须配置posthog/temporal/common/logger.py中的resolve_log_source函数能从 workflow 的 ID 和类型解析出log_source。设计哲学是日志在你需要时在场否则不挡路写 stdout 永远生效与 produce 前置条件是否满足无关即使 produce 条件不满足也不会让你的 workflow 崩溃。提示如果只关心 stdout可用get_write_only_logger获取只写 loggerget_produce_only_logger同理。get_logger应该只在模块顶部调用一次。如果会频繁打日志就在 activity/workflow 顶部对全局 logger 调用bind避免昂贵的全局查找bind还可以把变量直接绑定到 logger 实例上import dataclasses from structlog import get_logger from structlog.contextvars import bind_contextvars from temporalio import activity, workflow # Loggers initialized only once in module scope. LOGGER get_logger() # Takes a name for the logger, leave empty for module name activity.defn async def my_activity_with_lots_of_logging(inputs: MyActivityInputs): # All loggers from here on will have user_id in the context. # The context is copied over to new tasks and threads. bind_contextvars(user_idinputs.user_id) # We bind() as we will be logging a lot, and global lookups are expensive! logger LOGGER.bind() # Variables can be bound to the logger itself, these would only available to this logger: # logger LOGGER.bind(something..., another_one...) try: ... except: # Help your fellow PostHog engineers figure out what happened! logger.exception(Activity failed, reasonwow much technical) activity.defn async def my_small_activity(): # Using the global LOGGER has an associated performance cost. # But if you are not logging much, this is acceptable. LOGGER.info(Finished doing small thing) workflow.defn class MyWorfklow: workflow.run async def run(self): # Loggers can also be used in workflows logger LOGGER.bind() logger.info(Workflow start) await workflow.execute_activity(my_small_activity) await workflow.execute_activity(my_activity_with_lots_of_logging) logger.info(Workflow finished)开发 workflow/activity 时你大概率希望所有日志都带上 Temporal 上下文变量如activity_type、attempt、workflow_id。日志管道被配置为自动填充这些上下文——上一个示例中的所有日志都会带上这些变量。这条管道在本地同样生效无论是用phrocs还是手动运行start_temporal_worker.py启动本地 worker还是以如下方式运行单元测试DEBUG1 pytest path/to/your/tests.py -s本地日志由 structlog 用 Rich 渲染得到彩色、人类可读的输出而不是生产环境的 JSON 结构。警告logger 提供异步日志方法以 a 为前缀如await logger.ainfo。底层实现是开线程处理日志带来线程创建与上下文切换开销很多情况下并不划算更重要的是异步方法在 workflow 上下文中不可用因为 Temporal 运行 workflow 代码的自定义事件循环不实现任何线程调用。决策六留意 worker 关闭shutdownTemporal worker 在正常运维中经常被关停和重启例如新部署触发。关闭信号发出后worker 会等待其中正在运行的 activity 结束——但不会无限等待部署里配置了一个从几分钟到几小时不等的超时超时后所有 activity 被强杀。对长任务 activity 而言跟踪 worker 关闭时刻很有用activity 可以选择保存中间状态并提前退出避免被强杀而丢失全部进度。README 推荐posthog.temporal.common.shutdown中的ShutdownMonitor工具上下文管理器提供is_worker_shutdown()检查方法以及等待 worker 关闭的方法wait_for_worker_shutdown/wait_for_worker_shutdown_sync。短任务 activity 一般无需关心此事。从源码看ShutdownMonitor 对异步和同步两种环境都提供了支持start()起 asyncio taskstart_sync()起 daemon 线程并复制 context同时定义了WorkerShuttingDownError异常携带 activity/workflow/task queue 完整上下文——源码注释明确说明所有 Temporal activity 都应把WorkerShuttingDownError视为可重试异常这样新 worker 就能接手。raise_if_is_worker_shutdown()就是配合这一策略的便捷入口。而前面提到的Heartbeater内部也集成了关闭检测检测到 shutdown 时立即刷出最新心跳 details为跨 worker 的进度恢复争取时间。部署到生产workflow 和 activity 写完后还有几个步骤才能上生产。把 workflow 和 activity 分配到 workerPostHog 运行着多组 Temporal worker每组监听特定 task queue。哪组 worker 运行你的代码正是靠 task queue 协调的把 workflow/activity 放在某个 queue 上执行就只有配置为轮询该 queue 的那组 worker 会接走工作。由于每个产品对 worker 资源和行为的要求不同每个产品拥有自己的 worker 集合由产品团队管理其部署。对于还在开发中的 workflow/activity存在一组监听共享队列general-purpose-task-queue的 worker任何人都可以用它来跑开发期代码。一旦 workflow 走出原型阶段建议在 charts 仓库中找temporal-worker包创建自己的部署从而自定义资源限制、避免与共享队列上其他 workflow 的资源冲突。无论选哪个队列所有 worker 都在本仓库的代码里配置必须把你的 workflow 类和 activity 函数按所选队列加进 start_temporal_worker.py 中的映射WORKFLOWS_DICT/ACTIVITIES_DICT。推荐做法把产品的 workflow 和 activity 集中收集在产品包顶层__init__.py的两个列表里再导入到start_temporal_worker.py。假设你在products/travelling_salesman_solver/实现了产品其__init__.py形如from products.travelling_salesman_solver.temporal.workflows import MyWorkflow from products.travelling_salesman_solver.temporal.activity import my_first_activity, my_second_activity WORKFLOWS [MyWorkflow] ACTIVITIES [my_first_activity, my_second_activity]然后在start_temporal_worker.py中导入并加入字典from products.travelling_salesman_solver import WORKFLOWS as TS_WORKFLOWS, ACTIVITIES as TS_ACTIVITIES ... WORKFLOWS_DICT { ... # 若该产品已有独立部署可改用产品专属 task queue GENERAL_PURPOSE_TASK_QUEUE: TS_WORKFLOWS ... } ACTIVITIES_DICT { ... # 若该产品已有独立部署可改用产品专属 task queue GENERAL_PURPOSE_TASK_QUEUE: TS_ACTIVITIES ... }worker 部署后即可运行你的 workflow 和 activity。触发 worker 部署不可跳过这一步不是可选的后续事项charts 部署里的镜像指针charts 仓库的state/worker.yaml在 fleet 创建时手动播种此后只由本仓库container-images-cd.ymlworkflow 的 trigger 步骤更新。没有 trigger 步骤你的 fleet 会永远钉死在种子 digest 上——如果种子取自你的功能合入之前上线的 fleet 运行的就是不含你代码的镜像。worker 启动时如果 task queue 在镜像中没有注册任何 workflow会直接 crash-loop 并抛出ValueError: At least one activity, Nexus service, or workflow must be specified。务必在或随charts fleet PR 之前合入 trigger 步骤并确保种子 digest 的时间晚于功能合入提交。添加方式编辑container-images-cd.ymlGitHub workflow复制现有某个窄职责 worker 的 check trigger 步骤对release名必须与 charts state 文件的 key 一致。注意每个 trigger 步骤前都有一个 check 步骤确保只有特定模块的变更才触发 worker 重新部署而不是每次变更都触发。重启 worker 对正在运行的 workflow 有扰动所以应尽可能缩小触发重部署的变更范围。check 通常只需包含公共 temporal 模块 你产品专属模块——包括你的代码从仓库其他位置 import 的模块grep 一下入口的 import 列表。首次部署后用执行验证而不是用 dashboard 健康状态Temporal worker 部署禁用了 liveness/readiness 探针crash-loop 的 fleet 在 ArgoCD 里照样显示 HealthyPython 启动 traceback 走的是 info 级别的日志管道没有 error 级信号。正确姿势是查 fleet 日志service.name worker确认启动横幅每隔几分钟重复出现并确认一个真实 workflow 端到端跑完。注意Temporal worker 在关闭流程启动后会停止轮询新任务部署可配置关闭超时让正在运行的 workflow 有时间完成。正确配置能显著降低新部署的扰动但很难找到一个既保证所有任务跑完、又不至于永远等待的超时值。执行 workflow最直接的方式是execute_temporal_workflow命令对应 execute_temporal_workflow.pypython manage.py execute_temporal_workflow no-op {arg: 1, batch_export_id: test, team_id: 2}第一个参数是 workflow 名称第二个是 JSON 格式的 workflow 参数。重要该命令会向 Temporal service 发起请求。本地开发栈自带 Temporal service 并自动为该命令配置好但在生产执行 workflow 需要能访问生产 Temporal service 的服务器。警告注意你所处环境配置的TEMPORAL_TASK_QUEUE——如果该队列的 worker 没有注册你要执行的 workflow命令会失败。还有一种常见场景workflow 需要按固定间隔周期性执行而不是手动敲命令——这就是 Temporal schedule 的用武之地见上文 Schedule 四要素及 common/schedule.py 的封装。本地开发 Temporal开发栈包含一个充当本地编排器的 Temporal service、Temporal UI以及如果你使用phrocs多个自动启动的 Temporal worker——每个 task queue 一个。这些 worker 中包括监听共享general-purpose-task-queue的那个可直接用于开发。如果你新部署了一组 worker把它加进 bin/mprocs.yaml 让本地开发者也能自动拉起对应 worker。默认情况下Temporal worker 在 Python 文件变化时自动热重载与 backend 和 celery worker 类似。需要禁用热重载时设置TEMPORAL_DISABLE_HOT_RELOAD1。某些产品/特性可能需要额外配置例如数据仓库 workflow 需要环境变量中存在额外凭据可联系#team-data-warehouse团队获取帮助。请查阅你所开发产品的其他文档。运行 workflow 时日志会出现在 worker 日志里同时可以打开 Temporal UIhttp://localhost:8081查看 workflow 状态。测试模式Temporal 测试很贵要精打细算Temporal 测试开销很大——它们启动 Temporal 测试服务器、注册 activity而且经常被迫使用慢速的TransactionTestCase隔离模式。代价是真实的跨多个库的每测试 TRUNCATE所以要有意识地写。选最轻的、能证明目标的那层测试框架Harness适用场景成本纯 pytest无 Worker、无ActivityEnvironment单元测试纯函数——prompt 构建器、解析器、决策树~毫秒/测试ActivityEnvironment隔离测试单个 activity 函数体无 workflow 编排~几十毫秒/测试真实 Worker WorkflowEnvironment.start_time_skipping()端到端验证 workflow ↔ activity 编排的集成测试秒/测试 Worker 启动按 Temporal 官方指引绝大多数测试应写成仍能提供你所关心契约证明的最便宜 harness。拉起 Worker 是集成测试的刻意为之而不是默认选项。django_db(transactionTrue)规则通过sync_to_async/database_sync_to_async触碰 Django ORM 的异步测试不加pytest.mark.django_db(transactionTrue)会失败。原因asgiref 把被包装的同步函数派发到独立线程线程拿到自己的数据库连接看不到测试外层 atomic 包裹的事务。典型症状测试创建了行activity 却查不到。transactionTrue用每测试 TRUNCATE 全部替换了快速事务回滚。这是正确解法没有干净的办法在跨线程间共享 Django 连接但很贵所以只在需要时使用# 异步测试 Django ORM 需要 transactionTrue pytest.mark.django_db(transactionTrue) pytest.mark.asyncio async def test_my_activity(team): source await sync_to_async(ExternalDataSource.objects.create)(teamteam, ...) result await my_activity(SomeInputs(source_idsource.pk)) assert result.ok # 同步测试或不用 ORM 不需要 transactionTrue pytest.mark.django_db # transaction 回滚——快得多 def test_my_pure_function(team): source ExternalDataSource.objects.create(teamteam, ...) assert build_query(source) ...跨测试共享WorkflowEnvironmentWorker诱人但脆弱WorkflowEnvironment.start_time_skipping()会拉起真实的temporal-test-server进程Worker(...)注册全部 workflow activity。逐测试启动很贵于是很自然的优化是写一个 module 级 autouse fixture 启动一次、共享 client。PostHog 在test_end_to_end.py上尝试过并回退了PR #59405。对多数测试的节省是真实的约 8 秒/次但特定测试与进程级状态交互——mockShutdownMonitor.raise_if_is_worker_shutdown、patchee.api.billing.requests.get与posthog.cloud_utils.is_instance_licensed_cached——这些 mock 一旦作用到共享 Worker 的 activity 执行线程上测试 mock 上下文退出后要几分钟才能恢复。worker-shutdown 测试是最明显的受害者随后 billing-limits 测试也中招而且这种模式大概率不穷尽。CI 成本回归超过了节省。如果想再尝试共享安全前置条件是文件内没有任何测试会 mock/patch 共享 Worker 线程在 workflow 执行期间能触及的目标。这很难靠肉眼审查确认所以逃生舱口必不可少。该模式原理上可行但需要逐文件判断而不是全局优化。警惕生产代码中的connection.connect()调用某些 activity 会显式调用django.db.connection.connect()以从生产中过期的 worker 连接里恢复。但测试环境下这会摧毁测试的外层 atomic——随后在新连接上的下一次 ORM 读看不到任何数据teardown 会失败并报TransactionManagementError: The rollback flag doesnt work outside of an atomic block。定向修法在测试文件或整个包的conftest.py里把测试期间的 reconnect 置空操作pytest.fixture(autouseTrue) def _no_worker_reconnect(monkeypatch): from django.db import connection monkeypatch.setattr(connection, connect, lambda: None)加上这个 fixture 后测试就能留在快速pytest.mark.django_db模式无需transactionTrue。真实例子见 test_send_proxy_email.py。不要复制粘贴要参数化给新 source 变体加测试时每个端点一个测试很有诱惑力。这种诱惑的产物就是历史上test_end_to_end.py的状态83 个机械相同、只有 fixture 名不同的 10 行测试每个都支付全额的单测试成本。优先用pytest.mark.parametrizepytest.mark.django_db(transactionTrue) pytest.mark.asyncio pytest.mark.parametrize( schema_name,table_name,fixture_name, [ (charge, stripe_charge, stripe_charge), (customer, stripe_customer, stripe_customer), # ... 等等——新增端点现在只需追加一个元组 ], ) async def test_stripe_source(team, mock_stripe_client, request, schema_name, table_name, fixture_name): fixture_data request.getfixturevalue(fixture_name) await _run(teamteam, schema_nameschema_name, table_nametable_name, source_typeStripe, job_inputs_STRIPE_JOB_INPUTS, mock_data_responsefixture_data[data])每个用例的成本与 N 个独立函数相同但文件大幅瘦身新增端点也不再像一次重构。相关文档与 PostHog 中的实例Temporal 官方资料Python SDK 文档、schedules 文档、同步 vs 异步 activity 说明、测试最佳实践、SDK 仓库与示例仓库可作为延伸阅读。仓库内的实例整个 batch exports 都构建在 Temporal 之上各目的地S3、BigQuery、Snowflake、Databricks、Postgres、Redshift、Azure Blob、HTTP、文件下载等的 workflow 见 products/batch_exports/backend/temporal/destinationsTemporal workflow 的单元测试示例在 products/batch_exports/backend/tests/temporalDuckLake 数据建模写入同样利用 Temporal环境变量的 bucket 布局与 IAM 权限见 products/managed_warehouse/backend/README.md。posthog/temporal/目录本身还包含按领域划分的 workflow/activity 包——如ai/、ai_observability/、alerts/、data_modeling/、exports/、experiments/、session_replay/、proxy_service/、quota_limiting/等——它们全部经由 start_temporal_worker.py 中按 task queue 组织的WORKFLOWS_DICT/ACTIVITIES_DICT注册到对应 worker可作为新产品的直接参照。【免费下载链接】posthog:hedgehog: PostHog is the leading platform for building self-driving products. Our developer tools – AI observability, analytics, session replay, flags, experiments, error tracking, logs, and more – capture all the context agents need to diagnose problems, uncover opportunities, and ship fixes. Steer it all from Slack, web, desktop, or the MCP.项目地址: https://gitcode.com/GitHub_Trending/po/posthog创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考