把生产者和消费者连成父子后,这条链路再也结束不了
目标
拿一个只有一个文件的小队列,先通过转储文件看到用父子关系连接生产与消费会出什么问题,再改成链接。依次加上批量消费、队列等待时间、跨度种类、链接属性,再应用到扇出,最后顺着链接把一个订单的旅程跨越跟踪拼接起来。
为什么重要
父子关系的意思是“父级等待子级”。可是生产者不会等待消费结束。把这两者用父子关系连起来,用户明明已经得到响应,根跨度却关不掉,队列积压的日子里,一条跟踪会一连开着好几分钟。混进扇出之后,一个订单会长成数千个跨度,后端开始对这一条跟踪特殊对待。链接正是为这种位置而设——保留因果,但不等待。把消费一侧立为新跟踪的根,再用链接指向生产跨度,每条跟踪就能保持很小,整个旅程则可以顺着链接重新拼接起来。这与读取头部来接上断开的链条、或者在进程内传递上下文,是不同的判断。
步骤
- 创建
/root/tp-links/naive.py。转储路径优先读取TRACELAB_OUT,没有则使用/root/tp-links/01-naive.jsonl。队列文件放在与转储文件相同的目录中,命名为01-q.jsonl,开始时清空。在一个order.submit跨度里,用bus.publish放入 3 条消息(m-1–m-3),并用time.sleep(0.25)在队列中等待,然后在同一个跨度里用bus.poll取出,为每条消息创建子跨度order.handle并调用bus.handle。接着在/root/tp-links/01-problem.txt中写四行——traces=后面写转储文件中的跟踪数,root_span=后面写根跨度名称,wait_ms=后面写根跨度长度减去子级所覆盖区间之后的值(取整数),problem=后面写这种形态的问题是什么(至少 40 个字)。 - 创建
/root/tp-links/linked.py(默认转储路径/root/tp-links/02-linked.jsonl,队列文件02-q.jsonl)。在order.submit内为每条消息创建跨度order.publish,并把该跨度的trace_id和span_id以十六进制字符串随消息送出。等待之后取出消息的一方,要在order.submit之外运行,使跨度order.process成为跟踪的根,并用随消息而来的 ID 创建Link,通过links=挂上。两个跨度都在messaging.message.id属性中写下消息 ID。 - 创建
/root/tp-links/batch.py(默认转储路径/root/tp-links/03-batch.jsonl,队列文件03-q.jsonl)。这次放入 5 条消息(m-1–m-5),等待之后用bus.poll(Q, 5)一次性取出。用一个跨度order.process.batch处理取出的整批消息,在这个跨度上按取出的消息数挂上链接,并在整数属性messaging.batch.message_count中写下条数。这个跨度也必须是跟踪的根。 - 创建
/root/tp-links/wait.py(默认转储路径/root/tp-links/04-wait.jsonl,队列文件04-q.jsonl)。回到第 2 步的形态,但给消息增加produced_at_ns字段,随消息送出bus.now_ns()的值,并在消费跨度order.process上,把该时刻与现在的差值换算成毫秒,写入属性messaging.queue.wait_ms。然后在/root/tp-links/04-wait.tsv中,为每条消息写一行<메시지 아이디><탭><기다린 밀리초 소수 첫째 자리>(占位符依次为消息 ID、制表符、等待的毫秒数,保留到小数点后第一位),共三行。 - 创建
/root/tp-links/kinds.py(默认转储路径/root/tp-links/05-kinds.jsonl,队列文件05-q.jsonl)。给第 4 步的流程标上跨度种类——order.publish是SpanKind.PRODUCER,order.process是SpanKind.CONSUMER,外层的order.submit保持不变。并在这两个消息跨度上加messaging.system(labbus)、messaging.destination.name(orders)、messaging.operation.type(生产为send,消费为process)、messaging.operation.name、messaging.message.id。然后在/root/tp-links/05-kinds.tsv中写三行——每行是<스팬 이름><탭><종류><탭><operation.type>(占位符依次为跨度名称、制表符、种类、制表符、operation.type),顺序是order.submit、order.publish、order.process,order.submit的第三列是-。 - 创建
/root/tp-links/linkattrs.py(默认转储路径/root/tp-links/06-linkattrs.jsonl,队列文件06-q.jsonl)。流程与第 5 步相同,但创建Link时把属性作为第二个参数一起传入——link.relation写queue.message,messaging.message.id写该消息的 ID,messaging.destination.name写orders。转储文件的每个links列里都必须有这三个属性。 - 创建
/root/tp-links/fanout.py(默认转储路径/root/tp-links/07-fanout.jsonl,队列文件07-q.jsonl)。沿用第 6 步的规则,但扩展流程——order.submit产生 3 条消息,在消费一侧处理订单A-1002的order.process跨度内部再产生两条(跨度名称invoice.publish和email.publish,消息 IDinv-2和eml-2)。稍等片刻后把这两条取出,分别用invoice.process和email.process处理,它们也是新跟踪的根,并用链接指向前面的生产跨度。转储文件里应有 6 条跟踪。 - 创建
/root/tp-links/journey.py。用python3 journey.py <덤프> <메시지아이디>(占位符依次为转储文件、消息 ID)运行时,从创建该消息的order.publish跨度出发,按跳顺着链接相连的跨度往下走,每行输出<홉><탭><trace_id><탭><스팬이름>(占位符依次为跳数、制表符、trace_id、制表符、跨度名称)。出发的跨度是第 0 跳,下一跳是通过链接指向上一跳的跨度或该跨度的后代的那些跨度,同一跳内按跨度名称升序输出。无处可走时停止。把这个程序在/root/tp-links/07-fanout.jsonl和m-2上运行的输出,保存到/root/tp-links/08-journey.tsv。
参考
- 工作目录是
/root/tp-links。如果不存在,先创建。 - 带埋点的程序必须用
/opt/otel-lab/bin/python <파일>(占位符为文件)运行。系统python3中没有 OpenTelemetry SDK。只读取转储文件的程序用系统python3运行。 - 材料是
/opt/app/tracelab/tp_links/bus.py(只有一个文件的队列),公共接线是/opt/app/tracelab/dump.py,读取转储文件的辅助模块是/opt/lab/checks/_tplib.py。队列文件每一步各用一个,开始时清空。 - 常见错误:把消费循环留在
order.submit块里,只加了链接。那样既有链接又有父级,跟踪仍然连成一条。要通过转储文件中的parent_id是否为空来确认。 - Traces (OpenTelemetry Concepts) · Tracing API 规范 · 消息跨度语义约定 · 消息属性注册表 · Python 埋点文档
用父子关系连接生产与消费,会出什么事
创建 /root/tp-links/naive.py。转储路径优先读取 TRACELAB_OUT,没有则使用 /root/tp-links/01-naive.jsonl。队列文件放在与转储文件相同的目录中,命名为 01-q.jsonl,开始时清空。在一个 order.submit 跨度里,用 bus.publish 放入 3 条消息(m-1–m-3),并用 time.sleep(0.25) 在队列中等待,然后在同一个跨度里用 bus.poll 取出,为每条消息创建子跨度 order.handle 并调用 bus.handle。接着在 /root/tp-links/01-problem.txt 中写四行——traces= 后面写转储文件中的跟踪数,root_span= 后面写根跨度名称,wait_ms= 后面写根跨度长度减去子级所覆盖区间之后的值(取整数),problem= 后面写这种形态的问题是什么(至少 40 个字)。
材料是 /opt/app/tracelab/tp_links/bus.py——里面有 publish、poll、handle、now_ns。没被覆盖的区间,用 /opt/lab/checks/_tplib.py 的 covered_ns(부모, 자식들)(占位符依次为父级、子级们)求出。带埋点的程序用 /opt/otel-lab/bin/python 运行,重新生成转储文件之前要先删除它。
把消费立为新跟踪的根,并用链接连接
创建 /root/tp-links/linked.py(默认转储路径 /root/tp-links/02-linked.jsonl,队列文件 02-q.jsonl)。在 order.submit 内为每条消息创建跨度 order.publish,并把该跨度的 trace_id 和 span_id 以十六进制字符串随消息送出。等待之后取出消息的一方,要在 order.submit 之外运行,使跨度 order.process 成为跟踪的根,并用随消息而来的 ID 创建 Link,通过 links= 挂上。两个跨度都在 messaging.message.id 属性中写下消息 ID。
用 SpanContext(trace_id=..., span_id=..., is_remote=True, trace_flags=TraceFlags(TraceFlags.SAMPLED)) 创建上下文,再用 Link(ctx) 包起来。十六进制字符串用 format(값, "032x") 和 format(값, "016x") 生成,还原时用 int(문자열, 16)(三处占位符依次为值、值、字符串)。如果消费循环仍留在 with order.submit 块里,就不会成为根,看看缩进。
批量消费用一个带多个链接的跨度
创建 /root/tp-links/batch.py(默认转储路径 /root/tp-links/03-batch.jsonl,队列文件 03-q.jsonl)。这次放入 5 条消息(m-1–m-5),等待之后用 bus.poll(Q, 5) 一次性取出。用一个跨度 order.process.batch 处理取出的整批消息,在这个跨度上按取出的消息数挂上链接,并在整数属性 messaging.batch.message_count 中写下条数。这个跨度也必须是跟踪的根。
links= 传入的是列表——为每条消息创建 Link,作为 list 传入。每条消息各建一个跨度,会让处理一批的这一件事散开;不加链接只建一个跨度,又不知道哪些消息属于那一批。避开这两者的形态,就是这一步的答案。
如何测量在队列中等待的时间
创建 /root/tp-links/wait.py(默认转储路径 /root/tp-links/04-wait.jsonl,队列文件 04-q.jsonl)。回到第 2 步的形态,但给消息增加 produced_at_ns 字段,随消息送出 bus.now_ns() 的值,并在消费跨度 order.process 上,把该时刻与现在的差值换算成毫秒,写入属性 messaging.queue.wait_ms。然后在 /root/tp-links/04-wait.tsv 中,为每条消息写一行 <메시지 아이디><탭><기다린 밀리초 소수 첫째 자리>(占位符依次为消息 ID、制表符、等待的毫秒数,保留到小数点后第一位),共三行。
处理是逐条依次进行的,所以后取出的消息等得更久——三个值不相同才是正常的。时刻要由生产者而不是队列来打。如果把消费者取出的瞬间当作起点,等待时间就永远是 0。
PRODUCER 和 CONSUMER 什么时候用
创建 /root/tp-links/kinds.py(默认转储路径 /root/tp-links/05-kinds.jsonl,队列文件 05-q.jsonl)。给第 4 步的流程标上跨度种类——order.publish 是 SpanKind.PRODUCER,order.process 是 SpanKind.CONSUMER,外层的 order.submit 保持不变。并在这两个消息跨度上加 messaging.system(labbus)、messaging.destination.name(orders)、messaging.operation.type(生产为 send,消费为 process)、messaging.operation.name、messaging.message.id。然后在 /root/tp-links/05-kinds.tsv 中写三行——每行是 <스팬 이름><탭><종류><탭><operation.type>(占位符依次为跨度名称、制表符、种类、制表符、operation.type),顺序是 order.submit、order.publish、order.process,order.submit 的第三列是 -。
消息语义约定的表按操作种类决定跨度种类——创建和发送是 PRODUCER,应用程序处理消息的位置是 CONSUMER。包住业务流程的跨度不是消息跨度,所以不改它的种类。它会原样印在转储文件的 kind 列里。
在链接上写明为什么连在一起
创建 /root/tp-links/linkattrs.py(默认转储路径 /root/tp-links/06-linkattrs.jsonl,队列文件 06-q.jsonl)。流程与第 5 步相同,但创建 Link 时把属性作为第二个参数一起传入——link.relation 写 queue.message,messaging.message.id 写该消息的 ID,messaging.destination.name 写 orders。转储文件的每个 links 列里都必须有这三个属性。
像 Link(ctx, {"키": "값"})(占位符依次为键、值)这样,第二个参数就是属性。只有链接的话,“连着”这个事实会保留,但为什么连着不会保留——同样的两个跨度,可能因为队列连在一起,也可能因为重新处理而连在一起,这两种情况对读的人来说是完全不同的故事。
连扇出一并应用
创建 /root/tp-links/fanout.py(默认转储路径 /root/tp-links/07-fanout.jsonl,队列文件 07-q.jsonl)。沿用第 6 步的规则,但扩展流程——order.submit 产生 3 条消息,在消费一侧处理订单 A-1002 的 order.process 跨度内部再产生两条(跨度名称 invoice.publish 和 email.publish,消息 ID inv-2 和 eml-2)。稍等片刻后把这两条取出,分别用 invoice.process 和 email.process 处理,它们也是新跟踪的根,并用链接指向前面的生产跨度。转储文件里应有 6 条跟踪。
第二跳的生产跨度是消费跨度的子级——这是同一条跟踪内的同步调用,所以用父子关系才对。只有穿过队列时才用链接跳过去。哪种关系用什么来连,就是这一步的全部,数一数跟踪的数量,就能马上知道有没有分对。
顺着链接,把旅程重新拼接起来
创建 /root/tp-links/journey.py。用 python3 journey.py <덤프> <메시지아이디>(占位符依次为转储文件、消息 ID)运行时,从创建该消息的 order.publish 跨度出发,按跳顺着链接相连的跨度往下走,每行输出 <홉><탭><trace_id><탭><스팬이름>(占位符依次为跳数、制表符、trace_id、制表符、跨度名称)。出发的跨度是第 0 跳,下一跳是通过链接指向上一跳的跨度或该跨度的后代的那些跨度,同一跳内按跨度名称升序输出。无处可走时停止。把这个程序在 /root/tp-links/07-fanout.jsonl 和 m-2 上运行的输出,保存到 /root/tp-links/08-journey.tsv。
要走到第二跳,必须看到后代——因为创建发票消息的跨度是 order.process 的子级,而不是 order.process 本身。评分器会用其他消息 ID 来运行你的程序,所以不能把 m-2 写死在代码里。读取转储文件不需要 otel。