
本文从统一任务抽象、调度骨架、场景实现、质量门禁四个层面拆解。
任何自动化任务都可以拆成四要素:
把这四要素抽象成统一模型,就能用一套调度器驱动所有场景。
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Awaitable, Callable
import uuid, time
class Scene(str, Enum):
content = "content"
data = "data"
code = "code"
ops = "ops"
notify = "notify"
@dataclass
class Task:
scene: Scene
payload: dict[str, Any]
constraints: dict[str, Any] = field(default_factory=dict)
task_id: str = field(default_factory=lambda: uuid.uuid4().hex[:8])
retries: int = 0
max_retries: int = 3
status: str = "pending"
result: dict[str, Any] | None = None
created_at: float = field(default_factory=time.time)统一模型的价值:日志、重试、限流、审核、成本统计只写一次,所有场景复用。
import asyncio, logging
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
log = logging.getLogger("autopipe")
Handler = Callable[[Task], Awaitable[dict]]
HANDLERS: dict[Scene, Handler] = {}
def register(scene: Scene):
def deco(fn: Handler):
HANDLERS[scene] = fn
return fn
return deco
async def run_task(task: Task) -> Task:
handler = HANDLERS.get(task.scene)
if not handler:
task.status = "failed"
task.result = {"error": f"no handler for {task.scene}"}
return task
while task.retries <= task.max_retries:
start = time.time()
try:
task.result = await handler(task)
task.status = "succeeded"
log.info("scene=%s id=%s ok cost=%.2fs",
task.scene, task.task_id, time.time() - start)
return task
except Exception as e:
task.retries += 1
if task.retries > task.max_retries:
task.status = "failed"
task.result = {"error": str(e)}
log.error("scene=%s id=%s failed err=%s", task.scene, task.task_id, e)
return task
wait = 2 ** task.retries
log.warning("scene=%s id=%s retry=%d wait=%ds",
task.scene, task.task_id, task.retries, wait)
await asyncio.sleep(wait)
async def run_batch(tasks: list[Task], concurrency: int = 3) -> list[Task]:
sem = asyncio.Semaphore(concurrency)
async def wrap(t: Task):
async with sem:
return await run_task(t)
return await asyncio.gather(*(wrap(t) for t in tasks))要点:注册式解耦、并发受控、指数退避、失败不中断批次、日志可追踪。
import os, json
from openai import AsyncOpenAI
client = AsyncOpenAI(
api_key=os.getenv("OPENAI_API_KEY"),
base_url=os.getenv("OPENAI_BASE_URL"),
)
@register(Scene.content)
async def gen_content(task: Task) -> dict:
resp = await client.chat.completions.create(
model=os.getenv("OPENAI_MODEL", "gpt-4o-mini"),
messages=[
{"role": "system",
"content": "只输出 JSON:{title, points}。不编造事实,不含侵权元素。"},
{"role": "user", "content": task.payload["topic"]},
],
response_format={"type": "json_object"},
temperature=0.6,
)
return json.loads(resp.choices[0].message.content)import statistics
@register(Scene.data)
async def clean_data(task: Task) -> dict:
rows = task.payload["rows"]
clean = []
for r in rows:
try:
amount = float(r["amount"])
except (KeyError, ValueError, TypeError):
continue
clean.append({"id": r.get("id"), "amount": amount})
amounts = [r["amount"] for r in clean]
return {
"count": len(amounts),
"total": round(sum(amounts), 2),
"avg": round(statistics.mean(amounts), 2) if amounts else 0,
"max": max(amounts) if amounts else 0,
}注意:数据涉及个人信息时必须脱敏,不把原始数据贴给外部模型。
@register(Scene.code)
async def gen_tests(task: Task) -> dict:
source = task.payload["source"]
resp = await client.chat.completions.create(
model=os.getenv("OPENAI_MODEL", "gpt-4o-mini"),
messages=[
{"role": "system",
"content": "为函数生成 pytest,覆盖正常、边界、异常、空输入。只输出代码。"},
{"role": "user", "content": source},
],
temperature=0.2,
)
return {"tests": resp.choices[0].message.content}生成结果必须人工审查:断言是否有效、异常分支是否遗漏、mock 是否失真。
BAD_WORDS = {"违法", "暴力", "色情", "歧视", "虚假"}
DANGEROUS_CMD = ("rm -rf", "curl | sh", "chmod 777", "sudo")
def audit_text(text: str) -> None:
if any(w in text for w in BAD_WORDS):
raise ValueError("审核未通过")
def audit_commands(cmds: list[str]) -> list[str]:
return [c for c in cmds if any(d in c for d in DANGEROUS_CMD)]
def audit_task(task: Task) -> None:
audit_text(json.dumps(task.result, ensure_ascii=False, default=str))审核要覆盖输入和输出。运维类命令必须走白名单和沙箱,人工确认后再上生产。
async def main():
tasks = [
Task(Scene.content, {"topic": "自动化运维趋势"}),
Task(Scene.data, {"rows": [{"id": 1, "amount": "12.5"},
{"id": 2, "amount": "bad"}]}),
Task(Scene.code, {"source": "def add(a,b): return a+b"}),
]
results = await run_batch(tasks, concurrency=2)
for t in results:
print(t.scene, t.status, t.result)
# asyncio.run(main())多场景自动化生产的专业性,不在于覆盖多少场景,而在于建立统一任务抽象和调度骨架,让场景像插件一样注册进来,共享重试、限流、日志、审核和成本统计能力。代码可以简单,但权限、审核、日志、门禁和合规不能省。先跑通内容生成和数据清洗两个场景,再逐步扩展到代码、运维和通知分发。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。