TT Lab
开始
学习 学习路径 课程

分布式链路断掉的地方

把生产者和消费者连成父子后,这条链路再也结束不了

在 TT Lab 中继续学习

目标

拿一个只有一个文件的小队列,先通过转储文件看到用父子关系连接生产与消费会出什么问题,再改成链接。依次加上批量消费、队列等待时间、跨度种类、链接属性,再应用到扇出,最后顺着链接把一个订单的旅程跨越跟踪拼接起来。

为什么重要

父子关系的意思是“父级等待子级”。可是生产者不会等待消费结束。把这两者用父子关系连起来,用户明明已经得到响应,根跨度却关不掉,队列积压的日子里,一条跟踪会一连开着好几分钟。混进扇出之后,一个订单会长成数千个跨度,后端开始对这一条跟踪特殊对待。链接正是为这种位置而设——保留因果,但不等待。把消费一侧立为新跟踪的根,再用链接指向生产跨度,每条跟踪就能保持很小,整个旅程则可以顺着链接重新拼接起来。这与读取头部来接上断开的链条、或者在进程内传递上下文,是不同的判断。

步骤

  1. 创建 /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 个字)。
  2. 创建 /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。
  3. 创建 /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 中写下条数。这个跨度也必须是跟踪的根。
  4. 创建 /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、制表符、等待的毫秒数,保留到小数点后第一位),共三行。
  5. 创建 /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 的第三列是 -。
  6. 创建 /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 列里都必须有这三个属性。
  7. 创建 /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 条跟踪。
  8. 创建 /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/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。