プロデューサーとコンシューマーを親子にしたらトレースが終わらなくなった
目標
ファイル1つでできた小さなキューを題材に、生産と消費を親子でつないだときに起きる問題をダンプで見たあと、リンクに変えます。バッチ消費・キューの待ち時間・スパンの種類・リンク属性を順に付け、ファンアウトまで適用したら、リンクをたどって、1つの注文の道のりをトレースの向こう側までつなぎ合わせます。
なぜ重要なのか
親子は「親が子を待つ」という意味です。ところがプロデューサーは、消費が終わるのを待ちません。その2つを親子でつなぐと、ユーザーにはすでに応答が返っているのにルートスパンが閉じられず、キューが詰まった日には、トレース1つが数分も開いたままになります。ファンアウトが混ざると、注文1件が数千スパンに育ち、バックエンドがそのトレースだけを特別扱いし始めます。リンクは、まさにこの場面のためにあります。因果は残すが、待ちはしないという関係です。消費側を新しいトレースのルートとして立て、生産スパンをリンクで指せば、トレース1つ1つは小さく保たれ、全体の道のりはリンクをたどりながら再びつなぎ合わせられます。ヘッダーを読んで途切れたチェーンをつなぐ作業や、プロセス内でコンテキストを渡す作業とは、別の判断です。
ステップ
/root/tp-links/naive.pyを作成してください。ダンプのパスはTRACELAB_OUTから先に読み、なければ/root/tp-links/01-naive.jsonlを使います。キューファイルはダンプと同じディレクトリの01-q.jsonlとし、開始時に空にします。order.submitスパン1つの中で、bus.publishでメッセージ3件(m-1–m-3)を入れ、time.sleep(0.25)でキューで待ったあと、同じスパンの中でbus.pollで取り出し、メッセージごとに子スパンorder.handleを作ってbus.handleを呼びます。そのあと/root/tp-links/01-problem.txtに4行を書いてください。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を16進文字列にしてメッセージに載せて送ります。待ったあとに取り出す側は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)で一度に取り出します。取り出したバッチ全体を1つのスパン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に、メッセージごとに1行ずつ<메시지 아이디><탭><기다린 밀리초 소수 첫째 자리>(プレースホルダーはメッセージID、タブ、待ったミリ秒の小数第1位までの値です)を書いてください(3行)。/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はそのままにします。そして2つのメッセージングスパンに、messaging.system(labbus)・messaging.destination.name(orders)・messaging.operation.type(生産はsend、消費はprocess)・messaging.operation.name・messaging.message.idを付けます。そのあと/root/tp-links/05-kinds.tsvに3行を書いてください。各行は<스팬 이름><탭><종류><탭><operation.type>(プレースホルダーはスパン名、タブ、種類です)で、順序はorder.submit、order.publish、order.processであり、order.submitの3つ目の欄は-です。/root/tp-links/linkattrs.pyを作成してください(デフォルトのダンプのパスは/root/tp-links/06-linkattrs.jsonl、キューファイルは06-q.jsonl)。ステップ5と同じ流れですが、Linkを作るときに2番目の引数で属性も一緒に渡します。link.relationにqueue.message、messaging.message.idにそのメッセージのID、messaging.destination.nameにordersを書きます。ダンプのlinksの欄ごとに、この3つの属性が入っている必要があります。/root/tp-links/fanout.pyを作成してください(デフォルトのダンプのパスは/root/tp-links/07-fanout.jsonl、キューファイルは07-q.jsonl)。ステップ6のルールをそのまま使いながら、流れを広げます。order.submitがメッセージ3件を作り、消費側で注文A-1002を処理するorder.processスパンの中で、さらに2件を作ります(スパン名invoice.publishとemail.publish、メッセージIDinv-2とeml-2)。少し待ってからその2件を取り出し、invoice.processとemail.processでそれぞれ処理します。これらも新しいトレースのルートであり、リンクで前の生産スパンを指します。ダンプにはトレースが6個入っている必要があります。/root/tp-links/journey.pyを作成してください。python3 journey.py <덤프> <메시지아이디>(プレースホルダーはダンプとメッセージIDです)で実行すると、そのメッセージを作ったorder.publishスパンから出発し、リンクでつながったスパンをホップ単位でたどりながら、1行ずつ<홉><탭><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(ファイル1つでできたキュー)で、共通の配線は/opt/app/tracelab/dump.py、ダンプ読み込みのヘルパーは/opt/lab/checks/_tplib.pyです。キューファイルはステップごとに別に用意し、開始時に空にします。 - よくある間違い: 消費ループを
order.submitブロックの中に置いたまま、リンクだけを付けてしまいます。そうするとリンクも親もあるため、トレースが依然として1つにつながったままになります。ダンプの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スパン1つの中で、bus.publishでメッセージ3件(m-1–m-3)を入れ、time.sleep(0.25)でキューで待ったあと、同じスパンの中でbus.pollで取り出し、メッセージごとに子スパンorder.handleを作ってbus.handleを呼びます。そのあと/root/tp-links/01-problem.txtに4行を書いてください。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を16進文字列にしてメッセージに載せて送ります。待ったあとに取り出す側は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)で包みます。16進文字列はformat(값, "032x")とformat(값, "016x")で作り、元に戻すときはint(문자열, 16)です(プレースホルダーは値と文字列です)。消費ループがwith order.submitブロックの中に残っているとルートにならないので、インデントを確認してください。
バッチ消費はリンクが複数付いたスパン1つで
/root/tp-links/batch.pyを作成してください(デフォルトのダンプのパスは/root/tp-links/03-batch.jsonl、キューファイルは03-q.jsonl)。今回はメッセージを5件(m-1–m-5)入れ、待ったあとbus.poll(Q, 5)で一度に取り出します。取り出したバッチ全体を1つのスパンorder.process.batchで処理し、そのスパンに取り出したメッセージの数だけリンクを付け、整数属性messaging.batch.message_countに件数を書きます。このスパンもトレースのルートである必要があります。
links=にはリストを渡します。メッセージごとにLinkを作って、リストとして渡してください。メッセージごとにスパンを作ると、バッチを処理した1回の仕事が散らばり、リンクなしでスパン1つだけを作ると、どのメッセージがそのバッチに入っていたかがわかりません。両方を避ける形が、このステップの答えです。
キューで待った時間をどう測るか
/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に、メッセージごとに1行ずつ<메시지 아이디><탭><기다린 밀리초 소수 첫째 자리>(プレースホルダーはメッセージID、タブ、待ったミリ秒の小数第1位までの値です)を書いてください(3行)。
処理は1件ずつ順に行うので、あとに取り出したメッセージほど長く待ちます。3つの値が同じでないのが正常です。時刻はキューではなくプロデューサーが記録する必要があります。コンシューマーが取り出した瞬間を起点にすると、待った時間が常に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はそのままにします。そして2つのメッセージングスパンに、messaging.system(labbus)・messaging.destination.name(orders)・messaging.operation.type(生産はsend、消費はprocess)・messaging.operation.name・messaging.message.idを付けます。そのあと/root/tp-links/05-kinds.tsvに3行を書いてください。各行は<스팬 이름><탭><종류><탭><operation.type>(プレースホルダーはスパン名、タブ、種類です)で、順序はorder.submit、order.publish、order.processであり、order.submitの3つ目の欄は-です。
メッセージングのセマンティックコンベンションの表は、操作の種類でスパンの種類を決めます。作成と送信はPRODUCER、アプリケーションがメッセージを処理する場面はCONSUMERです。業務の流れを包むスパンはメッセージングスパンではないので、種類を変えません。ダンプのkindの欄に、名前のまま記録されます。
リンクに、なぜつながっているかを書く
/root/tp-links/linkattrs.pyを作成してください(デフォルトのダンプのパスは/root/tp-links/06-linkattrs.jsonl、キューファイルは06-q.jsonl)。ステップ5と同じ流れですが、Linkを作るときに2番目の引数で属性も一緒に渡します。link.relationにqueue.message、messaging.message.idにそのメッセージのID、messaging.destination.nameにordersを書きます。ダンプのlinksの欄ごとに、この3つの属性が入っている必要があります。
Link(ctx, {"키": "값"})のように、2番目の引数が属性です(プレースホルダーはキーと値です)。リンクだけだと「つながっている」という事実は残りますが、なぜつながっているかは残りません。同じ2つのスパンが、キューのためにつながることも、再処理のためにつながることもありますが、その2つは読む人にとってまったく別の話です。
1つが複数を生み出すファンアウトまで適用する
/root/tp-links/fanout.pyを作成してください(デフォルトのダンプのパスは/root/tp-links/07-fanout.jsonl、キューファイルは07-q.jsonl)。ステップ6のルールをそのまま使いながら、流れを広げます。order.submitがメッセージ3件を作り、消費側で注文A-1002を処理するorder.processスパンの中で、さらに2件を作ります(スパン名invoice.publishとemail.publish、メッセージIDinv-2とeml-2)。少し待ってからその2件を取り出し、invoice.processとemail.processでそれぞれ処理します。これらも新しいトレースのルートであり、リンクで前の生産スパンを指します。ダンプにはトレースが6個入っている必要があります。
2ホップ目の生産スパンは、消費スパンの子です。同じトレース内の同期呼び出しなので、親子が正しいです。リンクで飛び越えるのは、キューを通るときだけです。どの関係を何でつなぐかがこのステップのすべてで、トレース数を数えれば、正しく分かれたかがすぐわかります。
リンクをたどって道のりをつなぎ直す
/root/tp-links/journey.pyを作成してください。python3 journey.py <덤프> <메시지아이디>(プレースホルダーはダンプとメッセージIDです)で実行すると、そのメッセージを作ったorder.publishスパンから出発し、リンクでつながったスパンをホップ単位でたどりながら、1行ずつ<홉><탭><trace_id><탭><스팬이름>(プレースホルダーはホップ、タブ、スパン名です)を出力します。出発スパンがホップ0で、次のホップは前のホップのスパンまたはその子孫をリンクで指すスパンたちであり、同じホップ内ではスパン名の昇順で出力します。これ以上たどる先がなければ止まります。このプログラムを/root/tp-links/07-fanout.jsonlとm-2で実行した出力を、/root/tp-links/08-journey.tsvに保存してください。
2ホップ目へ進むには、子孫まで見る必要があります。請求書メッセージを作ったスパンはorder.processの子であって、order.process自身ではないからです。採点ツールは作成したプログラムをほかのメッセージIDでも実行するので、m-2をコードに埋め込んではいけません。ダンプを読むのにotelは不要です。