digraph root { graph [_draw_="c 9 -#fffffe00 C 7 -#44475a P 4 0 0 0 674 1721 674 1721 0 ", applyCss=dark, arrowsize=0.6, bb="0,0,1721,674", bgcolor="#44475a", charset="UTF-8", clusterize=false, color="#797976", fontcolor="#f8f8f2", fontname=Agave, intention=full, margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, rankdir=TD, shape=box, splines=spline, style=rounded, truecolor=true, xdotversion=1.7 ]; node [arrowsize=0.6, bgcolor="#44475a", color="#797976", fontcolor="#f8f8f2", fontname=Agave, label="\N", margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, shape=box, splines=true, style=rounded, truecolor=true ]; edge [arrowsize=0.6, bgcolor="#44475a", color="#797976", fontcolor="#f8f8f2", fontname=Agave, fontsize=6.8, margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, shape=box, splines=true, style=rounded, truecolor=true ]; subgraph G_cc_0 { graph [applyCss=dark, arrowsize=0.6, bb="", bgcolor="#44475a", charset="UTF-8", cluster="", clusterize=false, color="#797976", fontcolor="#f8f8f2", fontname=Agave, intention=full, margin=4.3, overlap=voronoi, pack="", packmode="", pencolor="#797976", penwidth=0.8, rankdir=TD, shape=box, splines=spline, style=rounded, truecolor=true ]; node [arrowsize=0.6, bgcolor="#44475a", color="#797976", fontcolor="#f8f8f2", fontname=Agave, fontsize="", height="", href="", label="\N", margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, pos="", shape=box, splines=true, style=rounded, truecolor=true, width="" ]; edge [arrowhead="", arrowsize=0.6, arrowtail="", bgcolor="#44475a", color="#797976", constrained="", dir="", fontcolor="#f8f8f2", fontname=Agave, fontsize=6.8, href="", margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, pos="", reverse="", shape=box, splines=true, style=rounded, tooltip="", truecolor=true, type="" ]; subgraph cluster_pattern__orchestrator__conclusion { graph [_draw_="S 17 -setlinewidth(0.3) c 7 -#c0c0c0 C 7 -#44475a b 25 667 96 667 96 1045 96 1045 96 1051 96 1057 102 1057 108 1057 108 1057 362 \ 1057 362 1057 368 1051 374 1045 374 1045 374 667 374 667 374 661 374 655 368 655 362 655 362 655 108 655 108 655 102 661 96 667 \ 96 ", applyCss=dark, arrowsize=0.6, bb="655,96,1057,374", bgcolor="#44475a", charset="UTF-8", cluster=true, clusterize=false, color="#797976", fontcolor="#f8f8f2", fontname=Agave, intention=full, margin=0.15, overlap=voronoi, pack=true, packmode=array, pencolor=silver, penwidth=0.3, rankdir=TD, shape=box, splines=spline, style=rounded, truecolor=true ]; node [arrowsize=0.6, bgcolor="#44475a", color="#797976", fontcolor="#f8f8f2", fontname=Agave, fontsize="", height="", href="", label="\N", margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, pos="", shape=box, splines=true, style=rounded, truecolor=true, width="" ]; edge [arrowhead="", arrowsize=0.6, arrowtail="", bgcolor="#44475a", color="#797976", constrained="", dir="", fontcolor="#f8f8f2", fontname=Agave, fontsize=6.8, href="", margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, pos="", reverse="", shape=box, splines=true, style=rounded, tooltip="", truecolor=true, type="" ]; pattern__orchestrator__conclusion [_ldraw_="F 9.5 5 -Agave c 7 -#f8f8f2 t 1 T 731 366.4 -1 46 10 -conclusion F 9.5 5 -Agave c 7 -#f8f8f2 t 0 T 731 353.4 -1 181 43 -• we \ use redis. is redisJSON needed? no. F 9.5 5 -Agave c 7 -#f8f8f2 T 731 343.4 -1 145 35 -• main queue: queue_orchestrator F 9.5 \ 5 -Agave c 7 -#f8f8f2 T 731 333.4 -1 145 35 -• no single long-running program F 9.5 5 -Agave c 7 -#f8f8f2 T 731 323.4 -1 190 \ 45 -• program is executed in steps, steps are F 9.5 5 -Agave c 7 -#f8f8f2 T 747 313.4 -1 158 35 -a disconnected sequence of \ messages ", fontsize=9.5, height=0.88889, href="#pattern__orchestrator__conclusion", label=<
conclusion
• we use redis. is redisJSON needed? no.
• main queue: queue_orchestrator
• no single long-running program
• program is executed in steps, steps are
a disconnected sequence of messages
>, pos="826,342", shape=plain, width=2.6667]; pattern__orchestrator__conclusion___pattern__orchestrator__completion_message [_ldraw_="F 4.75 5 -Agave c 7 -#f8f8f2 t 1 T 657 154.2 -1 41 18 -completion message F 4.75 5 -Agave c 7 -#f8f8f2 t 0 T 657 147.2 -1 55 27 \ -• pure rabbit-mq message F 4.75 5 -Agave c 7 -#ff79c6 T 712 147.2 -1 32 16 - program:next  F 4.75 5 -Agave c 7 -#f8f8f2 T \ 657 141.2 -1 66 32 -• advances queue_orchestrator F 4.75 5 -Agave c 7 -#f8f8f2 T 657 136.2 -1 93 44 -• must identify which \ step of program is F 4.75 5 -Agave c 7 -#f8f8f2 T 686.5 131.2 -1 34 15 -saying: move on F 4.75 5 -Agave c 7 -#f8f8f2 T 677 126.2 \ -1 46 23 -• if a bag of tasks: F 4.75 5 -Agave c 7 -#f8f8f2 T 717 121.2 -1 97 46 -• decr counter with redis eval. fires only \ F 4.75 5 -Agave c 7 -#f8f8f2 T 747 116.2 -1 37 16 -when counter===0 F 4.75 5 -Agave c 7 -#f8f8f2 T 717 111.2 -1 64 31 -• may \ postprocess its output ", fontsize=4.75, height=0.68056, href="#pattern__orchestrator__conclusion___pattern__orchestrator__completion_message", label=<
completion message
• pure rabbit-mq message program:next 
• advances queue_orchestrator
• must identify which step of program is
saying: move on
• if a bag of tasks:
• decr counter with redis eval. fires only
when counter===0
• may postprocess its output
>, pos="735,134", shape=plain, width=2.2083]; pattern__orchestrator__conclusion -> pattern__orchestrator__conclusion___pattern__orchestrator__completion_message [_draw_="S 15 -setlinewidth(2) S 7 -tapered c 9 -#fffffe00 C 7 -#797976 P 42 807.29 298.63 804.46 292.1 801.5 285.3 798.44 278.25 795.29 \ 270.99 792.07 263.56 788.79 256.01 785.47 248.37 782.13 240.67 778.79 232.96 775.45 225.28 772.14 217.65 768.88 210.13 765.67 202.75 \ 762.54 195.54 759.51 188.55 756.58 181.81 753.78 175.36 751.13 169.24 748.63 163.49 746.31 158.15 744.49 158.95 746.87 164.27 749.44 \ 169.98 752.18 176.07 755.06 182.48 758.07 189.18 761.2 196.13 764.42 203.3 767.72 210.64 771.09 218.12 774.49 225.7 777.93 233.34 \ 781.38 241.01 784.82 248.66 788.24 256.26 791.61 263.77 794.93 271.15 798.18 278.36 801.33 285.37 804.37 292.14 807.29 298.63 ", _tdraw_="S 15 -setlinewidth(2) S 5 -solid c 7 -#797976 C 7 -#797976 P 4 807.33 298.73 813.42 302.61 812.18 309.71 806.1 305.84 ", arrowsize=1.0, arrowtail=diamond, constrained=false, dir=back, href="#pattern__orchestrator__conclusion👉 pattern__orchestrator__conclusion___pattern__orchestrator__completion_message", penwidth=2, pos="s,812.18,309.71 807.29,298.63 788.57,256.27 760.46,192.64 745.4,158.55", reverse=false, style=tapered, tooltip="*= consistsOf", type=consistsOf]; pattern__orchestrator__conclusion___pattern__orchestrator__state [_ldraw_="F 4.75 5 -Agave c 7 -#f8f8f2 t 1 T 834 168.2 -1 12 5 -state F 4.75 5 -Agave c 7 -#f8f8f2 t 0 T 834 160.2 -1 68 33 -• central \ state of the program F 4.75 5 -Agave c 7 -#f8f8f2 T 834 155.2 -1 55 27 -• who holds the program? F 4.75 5 -Agave c 7 -#f8f8f2 \ T 854 151.2 -1 46 23 -• one big redis json F 4.75 5 -Agave c 7 -#ff79c6 T 900 151.2 -1 28 14 - program:id  F 4.75 5 -Agave \ c 7 -#f8f8f2 T 928 151.2 -1 3 1 -? F 4.75 5 -Agave c 7 -#f8f8f2 T 854 145.2 -1 86 41 -• or flat scalar json namespaced like \ F 4.75 5 -Agave c 7 -#f8f8f2 T 887.5 140.2 -1 19 8 -folders: F 4.75 5 -Agave c 7 -#ff79c6 T 834 129.45 -1 41 19 - /program:id/plan \ F 4.75 5 -Agave c 7 -#ff79c6 T 834 124.7 -1 55 24 -/program:id/plan/step:0 F 4.75 5 -Agave c 7 -#ff79c6 T 834 119.95 -1 59 26 -/\ program:id/outout/step:0 F 4.75 5 -Agave c 7 -#ff79c6 T 834 115.2 -1 68 30 -/program:id/plan/step:1/bag:0 F 4.75 5 -Agave c 7 \ -#ff79c6 T 834 110.45 -1 73 33 -/program:id/output/step:1/bag:0  F 4.75 5 -Agave c 7 -#f8f8f2 T 834 103.2 -1 77 37 -• redis \ prefers small values well F 4.75 5 -Agave c 7 -#f8f8f2 T 837.5 98.2 -1 70 31 -granulated into namespaced keys ", fontsize=4.75, height=1.0556, href="#pattern__orchestrator__conclusion___pattern__orchestrator__state", label=<
state
• central state of the program
• who holds the program?
• one big redis json program:id ?
• or flat scalar json namespaced like
folders:

 /program:id/plan
/program:id/plan/step:0
/program:id/outout/step:0
/program:id/plan/step:1/bag:0
/program:id/output/step:1/bag:0 

• redis prefers small values well
granulated into namespaced keys
>, pos="887,134", shape=plain, width=1.5]; pattern__orchestrator__conclusion -> pattern__orchestrator__conclusion___pattern__orchestrator__state [_draw_="S 15 -setlinewidth(2) S 7 -tapered c 9 -#fffffe00 C 7 -#797976 P 42 838.8 297.76 840.53 292.08 842.33 286.19 844.18 280.1 846.08 \ 273.86 848.02 267.48 850 260.99 852 254.41 854.02 247.77 856.06 241.09 858.09 234.4 860.13 227.73 862.15 221.09 864.15 214.52 866.12 \ 208.04 868.06 201.68 869.95 195.45 871.79 189.39 873.58 183.52 875.3 177.87 876.95 172.45 875.03 171.89 873.47 177.32 871.83 183 \ 870.14 188.9 868.38 194.99 866.58 201.24 864.74 207.64 862.87 214.15 860.97 220.75 859.05 227.41 857.12 234.12 855.19 240.83 853.26 \ 247.54 851.34 254.21 849.44 260.82 847.56 267.34 845.71 273.75 843.91 280.02 842.15 286.13 840.44 292.05 838.8 297.76 ", _tdraw_="S 15 -setlinewidth(2) S 5 -solid c 7 -#797976 C 7 -#797976 P 4 838.67 298.2 840.8 305.09 835.26 309.71 833.13 302.82 ", arrowsize=1.0, arrowtail=diamond, constrained=false, dir=back, href="#pattern__orchestrator__conclusion👉 pattern__orchestrator__conclusion___pattern__orchestrator__state", penwidth=2, pos="s,835.26,309.71 838.8,297.76 849.82,260.57 865.54,207.48 875.99,172.17", reverse=false, style=tapered, tooltip="*= consistsOf", type=consistsOf]; pattern__orchestrator__conclusion___pattern__orchestrator__task_output [_ldraw_="F 4.75 5 -Agave c 7 -#f8f8f2 t 1 T 960 141.2 -1 25 11 -task output F 4.75 5 -Agave c 7 -#f8f8f2 t 0 T 960 134.2 -1 88 42 -• each \ task when finished, can set its F 4.75 5 -Agave c 7 -#f8f8f2 T 960 129.45 -1 21 9 -output to F 4.75 5 -Agave c 7 -#ff79c6 T 981 \ 129.45 -1 75 35 - /program:id/plan/step:1/{bag:0}  F 4.75 5 -Agave c 7 -#f8f8f2 T 960 124.2 -1 86 41 -• this is done by the \ consumer, on ACK ", fontsize=4.75, height=0.31944, href="#pattern__orchestrator__conclusion___pattern__orchestrator__task_output", label=<
task output
• each task when finished, can set its
output to /program:id/plan/step:1/{bag:0} 
• this is done by the consumer, on ACK
>, pos="1008,134", shape=plain, width=1.3611]; pattern__orchestrator__conclusion -> pattern__orchestrator__conclusion___pattern__orchestrator__task_output [_draw_="S 15 -setlinewidth(2) S 7 -tapered c 9 -#fffffe00 C 7 -#797976 P 42 861.83 300.45 868.47 293.01 875.45 285.19 882.71 277.04 890.21 \ 268.63 897.9 260.02 905.71 251.26 913.61 242.41 921.52 233.54 929.41 224.7 937.22 215.94 944.9 207.34 952.38 198.95 959.63 190.82 \ 966.59 183.02 973.2 175.61 979.42 168.65 985.18 162.18 990.44 156.28 995.15 151.01 999.25 146.41 997.75 145.09 993.7 149.72 989.04 \ 155.04 983.84 160.99 978.13 167.51 971.99 174.54 965.45 182.02 958.57 189.88 951.4 198.08 943.99 206.54 936.4 215.22 928.68 224.05 \ 920.87 232.96 913.04 241.91 905.24 250.84 897.51 259.67 889.9 268.36 882.48 276.84 875.3 285.06 868.39 292.95 861.83 300.45 ", _tdraw_="S 15 -setlinewidth(2) S 5 -solid c 7 -#797976 C 7 -#797976 P 4 861.59 300.72 860.61 307.87 853.64 309.71 854.62 302.57 ", arrowsize=1.0, arrowtail=diamond, constrained=false, dir=back, href="#pattern__orchestrator__conclusion👉 pattern__orchestrator__conclusion___pattern__orchestrator__task_output", penwidth=2, pos="s,853.64,309.71 861.83,300.45 904.59,252.05 973.46,174.1 998.5,145.75", reverse=false, style=tapered, tooltip="*= consistsOf", type=consistsOf]; } subgraph cluster_pattern__orchestrator__problem__wait_for_bag_of_tasks { graph [_draw_="S 17 -setlinewidth(0.3) c 7 -#c0c0c0 C 7 -#44475a b 25 1087 0 1087 0 1255 0 1255 0 1261 0 1267 6 1267 12 1267 12 1267 362 1267 362 \ 1267 368 1261 374 1255 374 1255 374 1087 374 1087 374 1081 374 1075 368 1075 362 1075 362 1075 12 1075 12 1075 6 1081 0 1087 0 ", applyCss=dark, arrowsize=0.6, bb="1075,0,1267,374", bgcolor="#44475a", charset="UTF-8", cluster=true, clusterize=false, color="#797976", fontcolor="#f8f8f2", fontname=Agave, intention=full, margin=0.15, overlap=voronoi, pack=true, packmode=array, pencolor=silver, penwidth=0.3, rankdir=TD, shape=box, splines=spline, style=rounded, truecolor=true ]; node [arrowsize=0.6, bgcolor="#44475a", color="#797976", fontcolor="#f8f8f2", fontname=Agave, fontsize="", height="", href="", label="\N", margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, pos="", shape=box, splines=true, style=rounded, truecolor=true, width="" ]; edge [arrowhead="", arrowsize=0.6, arrowtail="", bgcolor="#44475a", color="#797976", constrained="", dir="", fontcolor="#f8f8f2", fontname=Agave, fontsize=6.8, href="", margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, pos="", reverse="", shape=box, splines=true, style=rounded, tooltip="", truecolor=true, type="" ]; pattern__orchestrator__problem__wait_for_bag_of_tasks [_ldraw_="F 9.5 5 -Agave c 7 -#f8f8f2 t 1 T 1076 366.4 -1 136 30 -problem: wait for bag of tasks F 9.5 5 -Agave c 7 -#f8f8f2 t 0 T 1076 353.4 \ -1 167 40 -• its called Barrier of orchestration F 9.5 5 -Agave c 7 -#f8f8f2 T 1076 343.4 -1 73 19 -• or logical cut F 9.5 \ 5 -Agave c 7 -#f8f8f2 T 1076 333.4 -1 190 45 -• its how to wait all parallel tasks have F 9.5 5 -Agave c 7 -#f8f8f2 T 1076 323.4 \ -1 185 41 -finished, and we shall proceed with next F 9.5 5 -Agave c 7 -#f8f8f2 T 1159.5 313.4 -1 23 5 -tasks ", fontsize=9.5, height=0.88889, href="#pattern__orchestrator__problem__wait_for_bag_of_tasks", label=<
problem: wait for bag of tasks
• its called Barrier of orchestration
• or logical cut
• its how to wait all parallel tasks have
finished, and we shall proceed with next
tasks
>, pos="1171,342", shape=plain, width=2.6667]; pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__Redis_Pub__Sub_notification [_ldraw_="F 4.75 5 -Agave c 7 -#f8f8f2 t 1 T 1076 178.2 -1 59 26 -Redis Pub/Sub notification F 4.75 5 -Agave c 7 -#f8f8f2 t 0 T 1076 170.2 \ -1 37 19 -• Redis: (best). F 4.75 5 -Agave c 7 -#f8f8f2 T 1076 165.2 -1 16 10 -• Setup F 4.75 5 -Agave c 7 -#f8f8f2 T 1076 \ 158.2 -1 82 36 -When you publish the parallel tasks: F 4.75 5 -Agave c 7 -#ff79c6 T 1076 153.2 -1 41 19 - Set counter: SET F 4.75 \ 5 -Agave c 7 -#ff79c6 T 1076 148.2 -1 52 23 -barrier:: N F 4.75 5 -Agave c 7 -#ff79c6 T 1076 143.2 -1 59 26 -Workers, \ when done: DECR F 4.75 5 -Agave c 7 -#ff79c6 T 1076 138.2 -1 55 24 -barrier:: If F 4.75 5 -Agave c 7 -#ff79c6 T 1076 \ 133.2 -1 79 37 -result == 0 → PUBLISH barrier:done F 4.75 5 -Agave c 7 -#ff79c6 T 1076 128.2 -1 32 15 -\":\"  F 4.75 \ 5 -Agave c 7 -#f8f8f2 T 1076 123.2 -1 43 19 -Orchestrator waits: F 4.75 5 -Agave c 7 -#ff79c6 T 1076 118.2 -1 84 38 - redisSub.subscribe(\"\ barrier:done\"); F 4.75 5 -Agave c 7 -#ff79c6 T 1076 113.2 -1 52 23 -redisSub.on(\"message\", F 4.75 5 -Agave c 7 -#ff79c6 T 1076 \ 108.2 -1 77 35 -(channel, msg) => { const 【wf, F 4.75 5 -Agave c 7 -#ff79c6 T 1076 103.2 -1 59 28 -step】 = msg.split(\":\"); \ F 4.75 5 -Agave c 7 -#ff79c6 T 1076 98.2 -1 73 33 -continueWorkflow(wf, step); });  F 4.75 5 -Agave c 7 -#f8f8f2 T 1076 93.2 -1 \ 142 47 -✅ zero polling ✅ instant reaction ✅ very F 4.75 5 -Agave c 7 -#f8f8f2 T 1114 88.2 -1 66 29 -scalableThis is the cleanest. ", fontsize=4.75, height=1.3333, href="#pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__Redis_Pub__Sub_notification", label=<
Redis Pub/Sub notification
• Redis: (best).
• Setup
When you publish the parallel tasks:
 Set counter: SET
barrier:<wf>:<step> N
Workers, when done: DECR
barrier:<wf>:<step> If
result == 0 → PUBLISH barrier:done
"<wf>:<step>" 

Orchestrator waits:
 redisSub.subscribe("barrier:done");
redisSub.on("message",
(channel, msg) => { const 【wf,
step】 = msg.split(":");
continueWorkflow(wf, step); }); 

✅ zero polling ✅ instant reaction ✅ very
scalableThis is the cleanest.
>, pos="1147,134", shape=plain, width=2]; pattern__orchestrator__problem__wait_for_bag_of_tasks -> pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__Redis_Pub__Sub_notification [_draw_="S 15 -setlinewidth(2) S 7 -tapered c 9 -#fffffe00 C 7 -#797976 P 42 1165.9 297.46 1165.35 292.3 1164.78 286.96 1164.19 281.46 1163.59 \ 275.83 1162.97 270.07 1162.35 264.22 1161.71 258.28 1161.07 252.28 1160.42 246.24 1159.77 240.17 1159.12 234.09 1158.46 228.03 1157.81 \ 221.99 1157.17 216.01 1156.53 210.09 1155.9 204.26 1155.28 198.53 1154.67 192.93 1154.07 187.48 1153.49 182.18 1151.51 182.42 1152.18 \ 187.7 1152.87 193.15 1153.57 198.73 1154.29 204.45 1155.02 210.26 1155.76 216.17 1156.51 222.14 1157.27 228.16 1158.02 234.22 1158.78 \ 240.28 1159.54 246.34 1160.29 252.37 1161.04 258.36 1161.77 264.28 1162.5 270.13 1163.22 275.87 1163.92 281.49 1164.6 286.98 1165.26 \ 292.31 1165.9 297.46 ", _tdraw_="S 15 -setlinewidth(2) S 5 -solid c 7 -#797976 C 7 -#797976 P 4 1165.94 297.8 1170.64 303.27 1167.4 309.71 1162.7 304.24 ", arrowsize=1.0, arrowtail=diamond, constrained=false, dir=back, href="#pattern__orchestrator__problem__wait_for_bag_of_tasks👉 pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__\ Redis_Pub__Sub_notification", penwidth=2, pos="s,1167.4,309.71 1165.9,297.46 1162,263.7 1156.6,216.98 1152.5,182.3", reverse=false, style=tapered, tooltip="*= consistsOf", type=consistsOf]; pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier [_ldraw_="F 4.75 5 -Agave c 7 -#f8f8f2 t 1 T 1238 133.2 -1 28 12 -join barrier ", fontsize=4.75, height=0.097222, href="#pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier", label=<
join barrier
>, pos="1252,134", shape=plain, width=0.41667]; pattern__orchestrator__problem__wait_for_bag_of_tasks -> pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier [_draw_="S 15 -setlinewidth(2) S 7 -tapered c 9 -#fffffe00 C 7 -#797976 P 42 1187.8 298.36 1191.14 289.99 1194.66 281.18 1198.32 272.01 1202.1 \ 262.56 1205.95 252.91 1209.86 243.15 1213.78 233.34 1217.68 223.57 1221.54 213.92 1225.32 204.46 1228.99 195.29 1232.51 186.47 1235.86 \ 178.09 1239.01 170.22 1241.91 162.95 1244.55 156.35 1246.88 150.51 1248.87 145.5 1250.5 141.41 1251.73 138.3 1249.87 137.58 1248.67 \ 140.69 1247.09 144.8 1245.16 149.83 1242.89 155.7 1240.34 162.33 1237.52 169.63 1234.47 177.54 1231.21 185.96 1227.79 194.82 1224.23 \ 204.03 1220.56 213.53 1216.81 223.23 1213.02 233.04 1209.21 242.89 1205.42 252.71 1201.68 262.4 1198.01 271.89 1194.46 281.1 1191.04 \ 289.95 1187.8 298.36 ", _tdraw_="S 15 -setlinewidth(2) S 5 -solid c 7 -#797976 C 7 -#797976 P 4 1187.72 298.55 1189.23 305.61 1183.3 309.71 1181.79 302.66 ", arrowsize=1.0, arrowtail=diamond, constrained=false, dir=back, href="#pattern__orchestrator__problem__wait_for_bag_of_tasks👉 pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__\ join_barrier", penwidth=2, pos="s,1183.3,309.71 1187.8,298.36 1209.1,244.06 1244.1,155.17 1250.8,137.94", reverse=false, style=tapered, tooltip="*= consistsOf", type=consistsOf]; "pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier___pattern__orchestrator__problem__wait_\ for_bag_of_tasks___pattern__orchestrator__Barrier_Counter_in_Redis__recommended_" [_ldraw_="F 2.38 5 -Agave c 7 -#f8f8f2 t 1 T 1085 48.1 -1 55 36 -Barrier Counter in Redis recommended F 2.38 5 -Agave c 7 -#f8f8f2 t 0 T 1085 \ 45.73 -1 50 33 -Simple, robust, works great with F 2.38 5 -Agave c 7 -#f8f8f2 T 1095.5 43.35 -1 34 22 -RabbitMQ.How it works: F \ 2.38 5 -Agave c 7 -#f8f8f2 T 1085 38.1 -1 65 48 -• Before publishing N parallel tasks → SET F 2.38 5 -Agave c 7 -#f8f8f2 T \ 1093 35.1 -1 49 32 -barrier::pending = N F 2.38 5 -Agave c 7 -#f8f8f2 T 1085 32.1 -1 65 46 -• Each worker, after \ finishing, runs: DECR F 2.38 5 -Agave c 7 -#f8f8f2 T 1096 29.1 -1 43 28 -barrier::pending F 2.38 5 -Agave c 7 -#f8f8f2 \ T 1085 26.1 -1 64 45 -• When the value hits 0, the orchestrator F 2.38 5 -Agave c 7 -#f8f8f2 T 1099.5 23.1 -1 35 23 -continues \ to next step. F 2.38 5 -Agave c 7 -#f8f8f2 T 1085 18.1 -1 26 17 -Why this is best: F 2.38 5 -Agave c 7 -#f8f8f2 T 1085 13.1 -1 13 \ 11 -• Atomic F 2.38 5 -Agave c 7 -#f8f8f2 T 1085 10.1 -1 31 23 -• No race conditions F 2.38 5 -Agave c 7 -#f8f8f2 T 1085 7.1 \ -1 28 21 -• Survives crashes F 2.38 5 -Agave c 7 -#f8f8f2 T 1085 2.1 -1 43 28 -Very easy in Node.js + STOMP ", fontsize=2.375, height=0.69444, href="#pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier___pattern__orchestrator__problem__wait_\ for_bag_of_tasks___pattern__orchestrator__Barrier_Counter_in_Redis__recommended_", label=<
Barrier Counter in Redis recommended
Simple, robust, works great with
RabbitMQ.How it works:
• Before publishing N parallel tasks → SET
barrier:<workflowId>:pending = N
• Each worker, after finishing, runs: DECR
barrier:<workflowId>:pending
• When the value hits 0, the orchestrator
continues to next step.
Why this is best:
• Atomic
• No race conditions
• Survives crashes
Very easy in Node.js + STOMP
>, pos="1117,25", shape=plain, width=0.93056]; pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier -> "pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier___pattern__orchestrator__problem__wait_\ for_bag_of_tasks___pattern__orchestrator__Barrier_Counter_in_Redis__recommended_" [_draw_="S 15 -setlinewidth(2) S 7 -tapered c 9 -#fffffe00 C 7 -#797976 P 122 1247.3 118.72 1246.72 117.17 1246.11 115.58 1245.45 113.96 \ 1244.76 112.3 1244.02 110.61 1243.25 108.9 1242.44 107.17 1241.59 105.44 1240.7 103.7 1239.76 101.96 1238.79 100.22 1237.78 98.5 \ 1236.73 96.79 1235.64 95.11 1234.5 93.45 1233.33 91.83 1232.11 90.25 1230.86 88.71 1229.56 87.22 1228.22 85.79 1224.71 82.31 1221.32 \ 79.23 1218.04 76.51 1214.85 74.11 1211.72 72 1208.64 70.13 1205.6 68.48 1202.57 67 1199.53 65.67 1196.47 64.42 1193.37 63.25 1190.2 \ 62.09 1186.96 60.92 1183.63 59.7 1180.17 58.4 1176.58 56.96 1172.83 55.36 1168.9 53.56 1164.78 51.51 1160.43 49.19 1159.99 48.95 \ 1159.54 48.7 1159.09 48.44 1158.64 48.19 1158.18 47.94 1157.73 47.69 1157.27 47.43 1156.82 47.18 1156.36 46.92 1155.9 46.66 1155.44 \ 46.4 1154.98 46.15 1154.52 45.89 1154.06 45.63 1153.6 45.37 1153.14 45.11 1152.67 44.85 1152.21 44.59 1151.75 44.32 1151.29 44.06 \ 1150.31 45.81 1150.78 46.07 1151.25 46.32 1151.72 46.57 1152.18 46.83 1152.65 47.08 1153.12 47.33 1153.58 47.59 1154.05 47.84 1154.51 \ 48.09 1154.98 48.34 1155.44 48.59 1155.9 48.84 1156.36 49.09 1156.82 49.34 1157.28 49.58 1157.74 49.83 1158.2 50.08 1158.65 50.32 \ 1159.11 50.56 1159.57 50.81 1163.99 53.08 1168.19 55.08 1172.19 56.84 1176 58.4 1179.64 59.79 1183.14 61.05 1186.51 62.22 1189.76 \ 63.33 1192.93 64.44 1196.03 65.56 1199.08 66.74 1202.11 68.01 1205.12 69.42 1208.15 71.01 1211.22 72.8 1214.34 74.83 1217.53 77.16 \ 1220.83 79.8 1224.24 82.81 1227.78 86.21 1229.13 87.61 1230.44 89.06 1231.71 90.57 1232.95 92.12 1234.14 93.71 1235.29 95.34 1236.41 \ 97 1237.48 98.68 1238.51 100.38 1239.51 102.1 1240.47 103.82 1241.38 105.54 1242.26 107.26 1243.1 108.97 1243.9 110.66 1244.66 112.34 \ 1245.38 113.99 1246.06 115.6 1246.7 117.18 1247.3 118.72 ", _tdraw_="S 15 -setlinewidth(2) S 5 -solid c 7 -#797976 C 7 -#797976 P 4 1247.31 118.74 1253.04 123.12 1251.2 130.09 1245.47 125.71 ", arrowsize=1.0, arrowtail=diamond, constrained=false, dir=back, href="#pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier👉 pattern__orchestrator__problem__\ wait_for_bag_of_tasks___pattern__orchestrator__join_barrier___pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__\ Barrier_Counter_in_Redis__recommended_", penwidth=2, pos="s,1251.2,130.09 1247.3,118.72 1243.5,108.6 1237.1,95.23 1228,86 1204,61.642 1190,66.334 1160,50 1157,48.367 1153.9,46.658 1150.8,\ 44.937", reverse=false, style=tapered, tooltip="*= consistsOf", type=consistsOf]; "pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier___pattern__orchestrator__problem__wait_\ for_bag_of_tasks___pattern__orchestrator__Barrier_via__completion_messages__queue" [_ldraw_="F 2.38 5 -Agave c 7 -#f8f8f2 t 1 T 1170 44.1 -1 59 43 -Barrier via “completion messages” queue F 2.38 5 -Agave c 7 -#f8f8f2 \ t 0 T 1175 41.73 -1 49 32 -RabbitMQ-only solution. Pattern: F 2.38 5 -Agave c 7 -#f8f8f2 T 1170 38.1 -1 61 43 -• Each parallel \ task sends a completion F 2.38 5 -Agave c 7 -#f8f8f2 T 1170 35.73 -1 44 29 -message to a dedicated queue: F 2.38 5 -Agave c 7 -#\ ff79c6 T 1214 35.73 -1 52 36 - workflow..step..completed  F 2.38 5 -Agave c 7 -#f8f8f2 T 1170 33.1 -1 47 34 -• Orchestrator \ counts how many F 2.38 5 -Agave c 7 -#f8f8f2 T 1171.5 30.1 -1 44 29 -\"completed\" messages arrived. F 2.38 5 -Agave c 7 -#f8f8f2 \ T 1170 27.1 -1 47 36 -• When it reaches N → continue. F 2.38 5 -Agave c 7 -#f8f8f2 T 1170 22.1 -1 8 5 -Cons: F 2.38 5 -Agave \ c 7 -#f8f8f2 T 1170 17.1 -1 65 46 -• Requires the orchestrator to stay online F 2.38 5 -Agave c 7 -#f8f8f2 T 1173.5 14.1 -1 \ 58 38 -and track counts in memory or storage. F 2.38 5 -Agave c 7 -#f8f8f2 T 1170 11.1 -1 58 41 -• Not safe without persistent \ storage. F 2.38 5 -Agave c 7 -#f8f8f2 T 1170 8.1 -1 47 34 -• Use only when you cannot add F 2.38 5 -Agave c 7 -#f8f8f2 T 1182 \ 5.1 -1 23 15 -Redis/Postgres. ", fontsize=2.375, height=0.59722, href="#pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier___pattern__orchestrator__problem__wait_\ for_bag_of_tasks___pattern__orchestrator__Barrier_via__completion_messages__queue", label=<
Barrier via “completion messages” queue
RabbitMQ-only solution. Pattern:
• Each parallel task sends a completion
message to a dedicated queue: workflow.<id>.step.<n>.completed 
• Orchestrator counts how many
"completed" messages arrived.
• When it reaches N → continue.
Cons:
• Requires the orchestrator to stay online
and track counts in memory or storage.
• Not safe without persistent storage.
• Use only when you cannot add
Redis/Postgres.
>, pos="1218,25", shape=plain, width=1.3611]; pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier -> "pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier___pattern__orchestrator__problem__wait_\ for_bag_of_tasks___pattern__orchestrator__Barrier_via__completion_messages__queue" [_draw_="S 15 -setlinewidth(2) S 7 -tapered c 9 -#fffffe00 C 7 -#797976 P 42 1247.5 118.89 1246.64 116.02 1245.71 112.95 1244.73 109.71 1243.71 \ 106.32 1242.64 102.8 1241.54 99.16 1240.41 95.42 1239.25 91.6 1238.08 87.73 1236.89 83.81 1235.69 79.87 1234.49 75.93 1233.3 72.01 \ 1232.12 68.12 1230.95 64.28 1229.8 60.52 1228.69 56.85 1227.6 53.29 1226.55 49.86 1225.55 46.58 1223.65 47.19 1224.74 50.44 1225.87 \ 53.84 1227.05 57.37 1228.27 61.01 1229.51 64.74 1230.78 68.55 1232.07 72.4 1233.36 76.29 1234.66 80.2 1235.96 84.1 1237.25 87.99 \ 1238.53 91.83 1239.79 95.61 1241.02 99.32 1242.22 102.93 1243.38 106.43 1244.49 109.79 1245.55 113 1246.56 116.04 1247.5 118.89 ", _tdraw_="S 15 -setlinewidth(2) S 5 -solid c 7 -#797976 C 7 -#797976 P 4 1247.52 118.94 1253.16 123.42 1251.2 130.36 1245.55 125.88 ", arrowsize=1.0, arrowtail=diamond, constrained=false, dir=back, href="#pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__join_barrier👉 pattern__orchestrator__problem__\ wait_for_bag_of_tasks___pattern__orchestrator__join_barrier___pattern__orchestrator__problem__wait_for_bag_of_tasks___pattern__orchestrator__\ Barrier_via__completion_messages__queue", penwidth=2, pos="s,1251.2,130.36 1247.5,118.89 1241.7,100.49 1231.4,68.136 1224.6,46.885", reverse=false, style=tapered, tooltip="*= consistsOf", type=consistsOf]; } subgraph cluster_pattern__orchestrator__problem__parallel_programs_ { graph [_draw_="S 17 -setlinewidth(0.3) c 7 -#c0c0c0 C 7 -#44475a b 25 1297 116.5 1297 116.5 1709 116.5 1709 116.5 1715 116.5 1721 122.5 1721 128.5 \ 1721 128.5 1721 430 1721 430 1721 436 1715 442 1709 442 1709 442 1297 442 1297 442 1291 442 1285 436 1285 430 1285 430 1285 128.5 \ 1285 128.5 1285 122.5 1291 116.5 1297 116.5 ", applyCss=dark, arrowsize=0.6, bb="1285,116.5,1721,442", bgcolor="#44475a", charset="UTF-8", cluster=true, clusterize=false, color="#797976", fontcolor="#f8f8f2", fontname=Agave, intention=full, margin=0.15, overlap=voronoi, pack=true, packmode=array, pencolor=silver, penwidth=0.3, rankdir=TD, shape=box, splines=spline, style=rounded, truecolor=true ]; node [arrowsize=0.6, bgcolor="#44475a", color="#797976", fontcolor="#f8f8f2", fontname=Agave, fontsize="", height="", href="", label="\N", margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, pos="", shape=box, splines=true, style=rounded, truecolor=true, width="" ]; edge [arrowhead="", arrowsize=0.6, arrowtail="", bgcolor="#44475a", color="#797976", constrained="", dir="", fontcolor="#f8f8f2", fontname=Agave, fontsize=6.8, href="", margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, pos="", reverse="", shape=box, splines=true, style=rounded, tooltip="", truecolor=true, type="" ]; pattern__orchestrator__problem__parallel_programs_ [_ldraw_="F 9.5 5 -Agave c 7 -#f8f8f2 t 1 T 1287 434.4 -1 122 27 -problem: parallel programs? F 9.5 5 -Agave c 7 -#f8f8f2 t 0 T 1287 421.4 \ -1 181 43 -• can multiple programs run in the same F 9.5 5 -Agave c 7 -#f8f8f2 T 1357 411.4 -1 41 9 -mq-space? F 9.5 5 -Agave \ c 7 -#f8f8f2 T 1287 401.4 -1 185 44 -• will a slow program not hold up a fast F 9.5 5 -Agave c 7 -#f8f8f2 T 1370 391.4 -1 19 \ 4 -one? F 9.5 5 -Agave c 7 -#f8f8f2 T 1287 381.4 -1 176 42 -• possible solution: sharding into few F 9.5 5 -Agave c 7 -#f8f8f2 \ T 1359 371.4 -1 32 7 -queue's F 9.5 5 -Agave c 7 -#f8f8f2 T 1287 359.4 -1 41 9 -You said: F 9.5 5 -Agave c 7 -#f8f8f2 T 1287 347.4 \ -1 131 32 -• fast workflow: occurs often F 9.5 5 -Agave c 7 -#f8f8f2 T 1287 337.4 -1 136 33 -• slow workflow: occurs rarely \ F 9.5 5 -Agave c 7 -#f8f8f2 T 1287 327.4 -1 140 34 -• slow workflow takes long time F 9.5 5 -Agave c 7 -#f8f8f2 T 1287 315.4 \ -1 336 75 -You want maximum simplicity, so:✅ Choose Option A: multiple orchestrator F 9.5 5 -Agave c 7 -#f8f8f2 T 1287 305.4 \ -1 221 49 -queuesYou only need two:/queue/orchestrator.fast F 9.5 5 -Agave c 7 -#f8f8f2 T 1287 295.4 -1 433 96 -/queue/orchestrator.slowThen \ route based on workflow type:const destination = workflow.type === F 9.5 5 -Agave c 7 -#f8f8f2 T 1287 285.4 -1 172 38 -'fast' ? '/\ queue/orchestrator.fast' : F 9.5 5 -Agave c 7 -#f8f8f2 T 1287 275.4 -1 280 62 -'/queue/orchestrator.slow'Each queue has its own \ orchestrator F 9.5 5 -Agave c 7 -#f8f8f2 T 1287 265.4 -1 228 51 -consumer.✅ Fast workflows never wait behind slow F 9.5 5 -Agave \ c 7 -#f8f8f2 T 1287 255.4 -1 205 47 -ones ✅ No complexity ✅ No risk ✅ Perfect F 9.5 5 -Agave c 7 -#f8f8f2 T 1483 245.4 -1 \ 41 9 -isolation ", fontsize=9.5, height=2.7778, href="#pattern__orchestrator__problem__parallel_programs_", label=<
problem: parallel programs?
• can multiple programs run in the same
mq-space?
• will a slow program not hold up a fast
one?
• possible solution: sharding into few
queue's
You said:
• fast workflow: occurs often
• slow workflow: occurs rarely
• slow workflow takes long time
You want maximum simplicity, so:✅ Choose Option A: multiple orchestrator
queuesYou only need two:/queue/orchestrator.fast
/queue/orchestrator.slowThen route based on workflow type:const destination = workflow.type ===
'fast' ? '/queue/orchestrator.fast' :
'/queue/orchestrator.slow'Each queue has its own orchestrator
consumer.✅ Fast workflows never wait behind slow
ones ✅ No complexity ✅ No risk ✅ Perfect
isolation
>, pos="1503,342", shape=plain, width=6.0417]; pattern__orchestrator__problem__parallel_programs____pattern__orchestrator__sharding [_ldraw_="F 4.75 5 -Agave c 7 -#f8f8f2 t 1 T 1458 147.2 -1 19 8 -sharding F 4.75 5 -Agave c 7 -#f8f8f2 t 0 T 1458 142.45 -1 88 39 -Multiple \ orchestrator queues, each one F 4.75 5 -Agave c 7 -#f8f8f2 T 1458 137.7 -1 75 33 -handles many workflows, but each F 4.75 5 -Agave \ c 7 -#f8f8f2 T 1458 132.95 -1 82 36 -workflow stays in exactly one shard F 4.75 5 -Agave c 7 -#f8f8f2 T 1458 128.2 -1 91 40 -based \ on hash(workflowId)% N. Gives you F 4.75 5 -Agave c 7 -#f8f8f2 T 1458 123.45 -1 84 37 -isolation, parallelism, and ordering F \ 4.75 5 -Agave c 7 -#f8f8f2 T 1475 118.7 -1 57 25 -with bounded queue count. ", fontsize=4.75, height=0.48611, href="#pattern__orchestrator__problem__parallel_programs____pattern__orchestrator__sharding", label=<
sharding
Multiple orchestrator queues, each one
handles many workflows, but each
workflow stays in exactly one shard
based on hash(workflowId)% N. Gives you
isolation, parallelism, and ordering
with bounded queue count.
>, pos="1503,134", shape=plain, width=1.2917]; pattern__orchestrator__problem__parallel_programs_ -> pattern__orchestrator__problem__parallel_programs____pattern__orchestrator__sharding [_draw_="S 15 -setlinewidth(2) S 7 -tapered c 9 -#fffffe00 C 7 -#797976 P 42 1503 229.68 1503.06 225.09 1503.12 220.51 1503.18 215.97 1503.23 \ 211.47 1503.29 207.02 1503.35 202.62 1503.4 198.29 1503.46 194.03 1503.51 189.85 1503.56 185.76 1503.61 181.77 1503.66 177.88 1503.71 \ 174.11 1503.76 170.46 1503.8 166.94 1503.85 163.56 1503.89 160.32 1503.93 157.24 1503.96 154.33 1504 151.58 1502 151.58 1502.04 \ 154.33 1502.07 157.24 1502.11 160.32 1502.15 163.56 1502.2 166.94 1502.24 170.46 1502.29 174.11 1502.34 177.88 1502.39 181.77 1502.44 \ 185.76 1502.49 189.85 1502.54 194.03 1502.6 198.29 1502.65 202.62 1502.71 207.02 1502.77 211.47 1502.82 215.97 1502.88 220.51 1502.94 \ 225.09 1503 229.68 ", _tdraw_="S 15 -setlinewidth(2) S 5 -solid c 7 -#797976 C 7 -#797976 P 4 1503 229.92 1507 235.92 1503 241.92 1499 235.92 ", arrowsize=1.0, arrowtail=diamond, constrained=false, dir=back, href="#pattern__orchestrator__problem__parallel_programs_👉 pattern__orchestrator__problem__parallel_programs____pattern__orchestrator__\ sharding", penwidth=2, pos="s,1503,241.92 1503,229.68 1503,198.99 1503,169.29 1503,151.58", reverse=false, style=tapered, tooltip="*= consistsOf", type=consistsOf]; } subgraph cluster_pattern__orchestrator { graph [_draw_="S 17 -setlinewidth(0.3) c 7 -#c0c0c0 C 7 -#44475a b 25 12 122 12 122 639 122 639 122 645 122 651 128 651 134 651 134 651 662 651 \ 662 651 668 645 674 639 674 639 674 12 674 12 674 6 674 0 668 0 662 0 662 0 134 0 134 0 128 6 122 12 122 ", applyCss=dark, arrowsize=0.6, bb="0,122,651,674", bgcolor="#44475a", charset="UTF-8", cluster=true, clusterize=false, color="#797976", fontcolor="#f8f8f2", fontname=Agave, intention=full, margin=0.15, overlap=voronoi, pack=true, packmode=array, pencolor=silver, penwidth=0.3, rankdir=TD, shape=box, splines=spline, style=rounded, truecolor=true ]; node [arrowsize=0.6, bgcolor="#44475a", color="#797976", fontcolor="#f8f8f2", fontname=Agave, fontsize="", height="", href="", label="\N", margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, pos="", shape=box, splines=true, style=rounded, truecolor=true, width="" ]; edge [arrowhead="", arrowsize=0.6, arrowtail="", bgcolor="#44475a", color="#797976", constrained="", dir="", fontcolor="#f8f8f2", fontname=Agave, fontsize=6.8, href="", margin=4.3, overlap=voronoi, pencolor="#797976", penwidth=0.8, pos="", reverse="", shape=box, splines=true, style=rounded, tooltip="", truecolor=true, type="" ]; pattern__orchestrator [_ldraw_="F 19 5 -Agave c 7 -#f8f8f2 t 1 T 250 658.8 -1 205 21 -pattern: orchestrator F 19 5 -Agave c 7 -#f8f8f2 t 0 T 250 634.8 -1 147 18 \ -• uses 3 queues F 19 5 -Agave c 7 -#f8f8f2 T 250 613.8 -1 400 44 -• structure: serial, exclusive+parallel, F 19 5 -Agave \ c 7 -#f8f8f2 T 420.5 592.8 -1 59 6 -serial F 19 5 -Agave c 7 -#f8f8f2 T 250 571.8 -1 98 13 -• example: F 19 5 -Agave c 7 -#f8f8f2 \ T 307 551.8 -1 20 5 -•  F 19 5 -Agave c 7 -#ff79c6 T 327 551.8 -1 199 24 - 【 serial tasks 】  F 19 5 -Agave c 7 -#f8f8f2 \ T 307 530.8 -1 20 5 -•  F 19 5 -Agave c 7 -#ff79c6 T 327 530.8 -1 287 33 - 【 bag of parallel tasks 】  F 19 5 -Agave c 7 \ -#f8f8f2 T 307 509.8 -1 20 5 -•  F 19 5 -Agave c 7 -#ff79c6 T 327 509.8 -1 160 20 - 【 untiloop 】  ", fontsize=19, height=2.3889, href="#pattern__orchestrator", label=<
pattern: orchestrator
• uses 3 queues
• structure: serial, exclusive+parallel,
serial
• example:
•  【 serial tasks 】 
•  【 bag of parallel tasks 】 
•  【 untiloop 】 
>, pos="450,588", shape=plain, width=5.5833]; pattern__orchestrator__Recommended_minimal_clean_setup__3_queues_ [_ldraw_="F 9.5 5 -Agave c 7 -#f8f8f2 t 1 T 2 458.4 -1 181 40 -Recommended minimal clean setup 3 queues F 9.5 5 -Agave c 7 -#f8f8f2 t 0 T \ 2 445.4 -1 95 24 -• /queue/orchestrator F 9.5 5 -Agave c 7 -#f8f8f2 T 22 435.4 -1 167 40 -• 100% serial queue: prefetch = \ 1, 1 F 9.5 5 -Agave c 7 -#f8f8f2 T 76 425.4 -1 59 13 -consumer only F 9.5 5 -Agave c 7 -#f8f8f2 T 22 415.4 -1 181 43 -• State \ machine messages only (continue F 9.5 5 -Agave c 7 -#f8f8f2 T 22 405.4 -1 163 36 -workflow, step done, parallel done, F 9.5 5 \ -Agave c 7 -#f8f8f2 T 65 395.4 -1 95 21 -start workflow, etc.) F 9.5 5 -Agave c 7 -#f8f8f2 T 2 385.4 -1 95 24 -• /queue/serial-tasks \ F 9.5 5 -Agave c 7 -#f8f8f2 T 29 375.4 -1 167 40 -• 100% serial queue: prefetch = 1, 1 F 9.5 5 -Agave c 7 -#f8f8f2 T 83 365.4 \ -1 59 13 -consumer only F 9.5 5 -Agave c 7 -#f8f8f2 T 29 355.4 -1 154 37 -• tasks run sequentially by nature F 9.5 5 -Agave c \ 7 -#f8f8f2 T 2 345.4 -1 104 26 -• /queue/parallel-tasks F 9.5 5 -Agave c 7 -#f8f8f2 T 58 335.4 -1 73 19 -• many consumers \ F 9.5 5 -Agave c 7 -#f8f8f2 T 58 325.4 -1 109 27 -• tasks run concurrently F 9.5 5 -Agave c 7 -#f8f8f2 T 2 315.4 -1 163 39 -• each \ task updates its state after F 9.5 5 -Agave c 7 -#f8f8f2 T 65 305.4 -1 37 8 -finished F 9.5 5 -Agave c 7 -#f8f8f2 T 67 295.4 -1 \ 73 19 -• TODO: but how? F 9.5 5 -Agave c 7 -#f8f8f2 T 67 285.4 -1 91 23 -• completion message F 9.5 5 -Agave c 7 -#f8f8f2 \ T 2 273.4 -1 176 39 -This is the simplest clean design that F 9.5 5 -Agave c 7 -#f8f8f2 T 67 263.4 -1 46 10 -gives you: F 9.5 5 \ -Agave c 7 -#f8f8f2 T 2 251.4 -1 149 36 -• controlled sequential execution F 9.5 5 -Agave c 7 -#f8f8f2 T 2 241.4 -1 104 26 -• unlimited \ parallelism F 9.5 5 -Agave c 7 -#f8f8f2 T 2 231.4 -1 118 29 -• simple, readable routing F 9.5 5 -Agave c 7 -#f8f8f2 T 2 221.4 \ -1 104 26 -• non-interfering logic ", fontsize=9.5, height=3.4444, href="#pattern__orchestrator__Recommended_minimal_clean_setup__3_queues_", label=<
Recommended minimal clean setup 3 queues
• /queue/orchestrator
• 100% serial queue: prefetch = 1, 1
consumer only
• State machine messages only (continue
workflow, step done, parallel done,
start workflow, etc.)
• /queue/serial-tasks
• 100% serial queue: prefetch = 1, 1
consumer only
• tasks run sequentially by nature
• /queue/parallel-tasks
• many consumers
• tasks run concurrently
• each task updates its state after
finished
• TODO: but how?
• completion message
This is the simplest clean design that
gives you:
• controlled sequential execution
• unlimited parallelism
• simple, readable routing
• non-interfering logic
>, pos="102,342", shape=plain, width=2.8194]; pattern__orchestrator -> pattern__orchestrator__Recommended_minimal_clean_setup__3_queues_ [_draw_="S 17 -setlinewidth(0.8) S 6 -dashed c 7 -#797976 B 7 265.17 501.82 246.61 490.7 228.53 478.72 212 466 211.82 465.86 211.65 465.73 \ 211.47 465.59 ", _hdraw_="S 17 -setlinewidth(0.8) S 5 -solid c 7 -#797976 C 7 -#797976 P 9 203.54 459.22 214.15 461.97 207.44 462.35 211.34 465.48 211.34 \ 465.48 211.34 465.48 207.44 462.35 208.52 468.99 203.54 459.22 ", arrowhead=open, arrowsize=1.0, href="#pattern__orchestrator👉 pattern__orchestrator__Recommended_minimal_clean_setup__3_queues_", pos="e,203.54,459.22 265.17,501.82 246.61,490.7 228.53,478.72 212,466 211.82,465.86 211.65,465.73 211.47,465.59", reverse=false, style=dashed, tooltip="-> isAnInputTo", type=isAnInputTo]; pattern__orchestrator__Workflow_Orchestration [_ldraw_="F 9.5 5 -Agave c 7 -#f8f8f2 t 1 T 223 441.4 -1 100 22 -Workflow Orchestration F 9.5 5 -Agave c 7 -#f8f8f2 t 0 T 223 431.9 -1 163 \ 36 -RabbitMQ is only the transport; the F 9.5 5 -Agave c 7 -#f8f8f2 T 223 422.4 -1 149 33 -orchestration logic lives in the F \ 9.5 5 -Agave c 7 -#f8f8f2 T 223 412.9 -1 217 52 -program-json.This is a “Workflow Orchestration” F 9.5 5 -Agave c 7 -#f8f8f2 \ T 223 403.4 -1 176 39 -(a.k.a. Saga Orchestration) pattern, a F 9.5 5 -Agave c 7 -#f8f8f2 T 223 393.9 -1 176 39 -classic Saga orchestration \ where every F 9.5 5 -Agave c 7 -#f8f8f2 T 293 384.4 -1 77 21 -“step” is either: F 9.5 5 -Agave c 7 -#f8f8f2 T 223 372.4 -1 \ 113 28 -• a sequential action, or F 9.5 5 -Agave c 7 -#f8f8f2 T 223 362.4 -1 176 42 -• a parallel fan-out group with a join \ F 9.5 5 -Agave c 7 -#f8f8f2 T 292.5 352.4 -1 37 8 -barrier. F 9.5 5 -Agave c 7 -#f8f8f2 T 223 340.4 -1 158 35 -Orchestrated Workflow \ with Barrier F 9.5 5 -Agave c 7 -#f8f8f2 T 268 330.4 -1 68 15 -Synchronization F 9.5 5 -Agave c 7 -#f8f8f2 T 223 318.4 -1 172 41 \ -• Reads the workflow definition (your F 9.5 5 -Agave c 7 -#f8f8f2 T 293 308.4 -1 32 7 -array). F 9.5 5 -Agave c 7 -#f8f8f2 \ T 223 298.4 -1 73 19 -• For each step: F 9.5 5 -Agave c 7 -#f8f8f2 T 223 288.4 -1 176 46 -• If it’s a single task → publish \ to a F 9.5 5 -Agave c 7 -#f8f8f2 T 250 278.4 -1 122 27 -serial queue, wait for ack. F 9.5 5 -Agave c 7 -#f8f8f2 T 223 268.4 -1 \ 167 44 -• If it’s a parallel block → fan-out F 9.5 5 -Agave c 7 -#f8f8f2 T 223 258.4 -1 172 38 -messages to a parallel queue, \ track N F 9.5 5 -Agave c 7 -#f8f8f2 T 223 248.4 -1 136 30 -acks, wait until all are done F 9.5 5 -Agave c 7 -#f8f8f2 T 236.5 238.4 \ -1 145 32 -(barrier),continue to next step. ", fontsize=9.5, height=2.9722, href="#pattern__orchestrator__Workflow_Orchestration", label=<
Workflow Orchestration
RabbitMQ is only the transport; the
orchestration logic lives in the
program-json.This is a “Workflow Orchestration”
(a.k.a. Saga Orchestration) pattern, a
classic Saga orchestration where every
“step” is either:
• a sequential action, or
• a parallel fan-out group with a join
barrier.
Orchestrated Workflow with Barrier
Synchronization
• Reads the workflow definition (your
array).
• For each step:
• If it’s a single task → publish to a
serial queue, wait for ack.
• If it’s a parallel block → fan-out
messages to a parallel queue, track N
acks, wait until all are done
(barrier),continue to next step.
>, pos="331,342", shape=plain, width=3.0417]; pattern__orchestrator -> pattern__orchestrator__Workflow_Orchestration [_draw_="S 15 -setlinewidth(2) S 7 -tapered c 9 -#fffffe00 C 7 -#797976 P 42 403.1 490.84 402.15 488.77 401.19 486.69 400.23 484.61 399.27 \ 482.53 398.3 480.44 397.33 478.34 396.37 476.24 395.4 474.14 394.43 472.04 393.46 469.93 392.48 467.82 391.51 465.71 390.54 463.6 \ 389.57 461.5 388.6 459.39 387.62 457.28 386.65 455.18 385.68 453.08 384.72 450.98 383.75 448.88 381.95 449.76 383.01 451.81 384.07 \ 453.87 385.13 455.92 386.19 457.98 387.25 460.05 388.31 462.11 389.37 464.17 390.44 466.24 391.5 468.3 392.56 470.37 393.62 472.43 \ 394.68 474.49 395.74 476.55 396.8 478.6 397.86 480.65 398.91 482.7 399.96 484.74 401.01 486.78 402.06 488.81 403.1 490.84 ", _tdraw_="S 15 -setlinewidth(2) S 5 -solid c 7 -#797976 C 7 -#797976 P 4 403.12 490.88 409.34 494.51 408.38 501.66 402.15 498.02 ", arrowsize=1.0, arrowtail=diamond, constrained=false, dir=back, href="#pattern__orchestrator👉 pattern__orchestrator__Workflow_Orchestration", penwidth=2, pos="s,408.38,501.66 403.1,490.84 396.45,477.21 389.59,463.13 382.85,449.32", reverse=false, style=tapered, tooltip="*= consistsOf", type=consistsOf]; pattern__orchestrator__experimental [_ldraw_="F 9.5 5 -Agave c 7 -#f8f8f2 t 1 T 460 366.4 -1 55 12 -experimental F 9.5 5 -Agave c 7 -#f8f8f2 t 0 T 460 354.4 -1 140 34 -• TODO: \ could also support eval F 9.5 5 -Agave c 7 -#ff79c6 T 600 354.4 -1 28 8 - code  F 9.5 5 -Agave c 7 -#f8f8f2 T 628 354.4 -1 5 1 \ -? F 9.5 5 -Agave c 7 -#f8f8f2 T 460 343.4 -1 190 45 -• holds program state in namespaced redis F 9.5 5 -Agave c 7 -#f8f8f2 \ T 521 333.4 -1 68 15 -keys: flat json F 9.5 5 -Agave c 7 -#f8f8f2 T 460 323.4 -1 176 42 -• state can collect output of tasks \ in F 9.5 5 -Agave c 7 -#f8f8f2 T 507 313.4 -1 82 18 -adjacent structure ", fontsize=9.5, height=0.88889, href="#pattern__orchestrator__experimental", label=<
experimental
• TODO: could also support eval code ?
• holds program state in namespaced redis
keys: flat json
• state can collect output of tasks in
adjacent structure
>, pos="555,342", shape=plain, width=2.6667]; pattern__orchestrator -> pattern__orchestrator__experimental [_draw_="S 17 -setlinewidth(0.8) S 6 -dashed c 7 -#797976 B 4 486.72 501.66 503.98 461.56 523.75 415.63 537.55 383.54 ", _hdraw_="S 17 -setlinewidth(0.8) S 5 -solid c 7 -#797976 C 7 -#797976 P 9 541.64 374.05 541.81 385.01 539.66 378.64 537.68 383.23 537.68 \ 383.23 537.68 383.23 539.66 378.64 533.55 381.45 541.64 374.05 ", arrowhead=open, arrowsize=1.0, href="#pattern__orchestrator👉 pattern__orchestrator__experimental", pos="e,541.64,374.05 486.72,501.66 503.98,461.56 523.75,415.63 537.55,383.54", reverse=false, style=dashed, tooltip="-> isAnInputTo", type=isAnInputTo]; pattern__orchestrator__Recommended_minimal_clean_setup__3_queues____pattern__orchestrator__queues [_ldraw_="F 4.75 5 -Agave c 7 -#f8f8f2 t 1 T 77 142.2 -1 14 6 -queues F 4.75 5 -Agave c 7 -#f8f8f2 t 0 T 77 134.2 -1 46 23 -• queue_orchestrator \ F 4.75 5 -Agave c 7 -#f8f8f2 T 77 129.2 -1 46 23 -• queue_tasks_serial F 4.75 5 -Agave c 7 -#f8f8f2 T 77 124.2 -1 50 25 -• queue_\ tasks_parallel ", fontsize=4.75, height=0.33333, href="#pattern__orchestrator__Recommended_minimal_clean_setup__3_queues____pattern__orchestrator__queues", label=<
queues
• queue_orchestrator
• queue_tasks_serial
• queue_tasks_parallel
>, pos="102,134", shape=plain, width=0.72222]; pattern__orchestrator__Recommended_minimal_clean_setup__3_queues_ -> pattern__orchestrator__Recommended_minimal_clean_setup__3_queues____pattern__orchestrator__queues [_draw_="S 17 -setlinewidth(0.8) S 6 -dashed c 7 -#797976 B 4 102 217.9 102 194.21 102 171.99 102 156.53 ", _hdraw_="S 17 -setlinewidth(0.8) S 5 -solid c 7 -#797976 C 7 -#797976 P 9 102 146.12 106.5 156.12 102 151.12 102 156.12 102 156.12 102 156.12 \ 102 151.12 97.5 156.12 102 146.12 ", arrowhead=open, arrowsize=1.0, href="#pattern__orchestrator__Recommended_minimal_clean_setup__3_queues_👉 pattern__orchestrator__Recommended_minimal_clean_setup__3_\ queues____pattern__orchestrator__queues", pos="e,102,146.12 102,217.9 102,194.21 102,171.99 102,156.53", reverse=false, style=dashed, tooltip="-> isAnInputTo", type=isAnInputTo]; } pattern__orchestrator -> pattern__orchestrator__conclusion [_draw_="S 15 -setlinewidth(2) S 7 -tapered c 9 -#fffffe00 C 7 -#797976 P 82 608.16 495.59 610.56 494.13 612.95 492.67 615.35 491.21 617.74 \ 489.74 620.13 488.27 622.51 486.8 624.89 485.33 627.26 483.86 629.63 482.39 631.99 480.92 634.34 479.44 636.69 477.97 639.03 476.5 \ 641.36 475.03 643.68 473.55 645.99 472.08 648.3 470.61 650.59 469.15 652.87 467.68 655.14 466.22 662.02 461.74 668.95 457.18 675.92 \ 452.54 682.92 447.84 689.93 443.09 696.93 438.3 703.91 433.48 710.86 428.66 717.75 423.84 724.57 419.04 731.3 414.27 737.94 409.55 \ 744.46 404.88 750.84 400.28 757.08 395.77 763.16 391.36 769.06 387.05 774.76 382.88 780.25 378.84 785.52 374.95 784.32 373.35 779.09 \ 377.28 773.64 381.37 767.98 385.6 762.13 389.95 756.1 394.42 749.9 398.99 743.56 403.64 737.09 408.37 730.5 413.16 723.81 417.99 \ 717.04 422.85 710.2 427.73 703.3 432.61 696.37 437.49 689.42 442.34 682.46 447.16 675.5 451.92 668.58 456.62 661.69 461.25 654.86 \ 465.78 652.6 467.27 650.34 468.75 648.06 470.24 645.77 471.73 643.47 473.23 641.16 474.72 638.85 476.21 636.52 477.71 634.19 479.2 \ 631.85 480.7 629.5 482.19 627.15 483.68 624.79 485.18 622.43 486.67 620.06 488.16 617.68 489.65 615.31 491.14 612.93 492.63 610.54 \ 494.11 608.16 495.59 ", _tdraw_="S 15 -setlinewidth(2) S 5 -solid c 7 -#797976 C 7 -#797976 P 4 607.92 495.74 604.89 502.28 597.68 502 600.71 495.46 ", arrowsize=1.0, arrowtail=diamond, constrained=false, dir=back, href="#pattern__orchestrator👉 pattern__orchestrator__conclusion", penwidth=2, pos="s,597.68,502 608.16,495.59 624.11,485.8 639.95,475.82 655,466 700.47,436.31 750.73,399.68 784.92,374.15", reverse=false, style=tapered, tooltip="*= consistsOf", type=consistsOf]; pattern__orchestrator -> pattern__orchestrator__problem__wait_for_bag_of_tasks [_draw_="S 15 -setlinewidth(2) S 7 -tapered c 9 -#fffffe00 C 7 -#797976 P 82 663.27 578.15 681.95 576.1 701.05 573.78 720.51 571.17 740.3 \ 568.26 760.35 565.02 780.63 561.44 801.08 557.51 821.67 553.21 842.33 548.52 863.03 543.42 883.71 537.91 904.34 531.96 924.85 525.55 \ 945.2 518.68 965.36 511.32 985.25 503.46 1004.85 495.08 1024.11 486.16 1042.96 476.7 1061.38 466.66 1067.2 463.19 1072.91 459.45 \ 1078.52 455.48 1084.02 451.3 1089.4 446.94 1094.65 442.4 1099.78 437.73 1104.78 432.94 1109.64 428.05 1114.35 423.09 1118.92 418.09 \ 1123.33 413.06 1127.59 408.03 1131.68 403.03 1135.6 398.07 1139.36 393.18 1142.93 388.38 1146.32 383.7 1149.52 379.17 1152.53 374.8 \ 1150.87 373.66 1147.89 378.04 1144.72 382.57 1141.36 387.24 1137.82 392.02 1134.11 396.91 1130.22 401.86 1126.16 406.85 1121.95 \ 411.87 1117.57 416.89 1113.05 421.88 1108.38 426.83 1103.57 431.71 1098.62 436.49 1093.54 441.15 1088.34 445.67 1083.01 450.03 1077.58 \ 454.2 1072.03 458.15 1066.37 461.88 1060.62 465.34 1042.29 475.42 1023.51 484.93 1004.33 493.9 984.79 502.34 964.95 510.26 944.85 \ 517.68 924.54 524.62 904.08 531.1 883.49 537.12 862.84 542.7 842.18 547.87 821.54 552.64 800.98 557.01 780.55 561.02 760.29 564.67 \ 740.25 567.98 720.48 570.97 701.03 573.65 681.94 576.04 663.27 578.15 ", _tdraw_="S 15 -setlinewidth(2) S 5 -solid c 7 -#797976 C 7 -#797976 P 4 663.03 578.17 657.46 582.76 651.09 579.39 656.65 574.8 ", arrowsize=1.0, arrowtail=diamond, constrained=false, dir=back, href="#pattern__orchestrator👉 pattern__orchestrator__problem__wait_for_bag_of_tasks", penwidth=2, pos="s,651.09,579.39 663.27,578.15 786.31,565.16 940.1,534.96 1061,466 1099.9,443.79 1132.4,402.77 1151.7,374.23", reverse=false, style=tapered, tooltip="*= consistsOf", type=consistsOf]; pattern__orchestrator -> pattern__orchestrator__problem__parallel_programs_ [_draw_="S 15 -setlinewidth(2) S 7 -tapered c 9 -#fffffe00 C 7 -#797976 P 82 663.36 575.52 689.15 573.37 715.91 570.95 743.58 568.27 772.06 \ 565.29 801.27 562.02 831.13 558.43 861.56 554.51 892.47 550.24 923.78 545.62 955.42 540.62 987.29 535.23 1019.32 529.44 1051.42 \ 523.23 1083.52 516.59 1115.52 509.5 1147.35 501.95 1178.92 493.93 1210.16 485.42 1240.97 476.4 1271.28 466.86 1274.31 465.86 1277.34 \ 464.83 1280.37 463.78 1283.41 462.71 1286.46 461.62 1289.5 460.5 1292.55 459.36 1295.6 458.21 1298.65 457.03 1301.7 455.83 1304.75 \ 454.62 1307.8 453.38 1310.85 452.13 1313.9 450.87 1316.95 449.58 1319.99 448.28 1323.03 446.97 1326.06 445.64 1329.09 444.3 1332.11 \ 442.94 1331.29 441.12 1328.28 442.48 1325.26 443.83 1322.24 445.16 1319.21 446.48 1316.18 447.79 1313.15 449.08 1310.11 450.35 1307.08 \ 451.61 1304.04 452.84 1301 454.06 1297.96 455.27 1294.92 456.45 1291.89 457.61 1288.85 458.75 1285.82 459.87 1282.79 460.97 1279.77 \ 462.05 1276.74 463.1 1273.73 464.13 1270.72 465.14 1240.47 474.76 1209.71 483.85 1178.53 492.45 1147 500.56 1115.22 508.2 1083.26 \ 515.37 1051.2 522.11 1019.13 528.41 987.13 534.29 955.28 539.77 923.67 544.86 892.37 549.58 861.48 553.93 831.07 557.94 801.22 561.62 \ 772.02 564.98 743.55 568.03 715.9 570.8 689.14 573.29 663.36 575.52 ", _tdraw_="S 15 -setlinewidth(2) S 5 -solid c 7 -#797976 C 7 -#797976 P 4 663.32 575.52 657.66 579.99 651.36 576.48 657.02 572.01 ", arrowsize=1.0, arrowtail=diamond, constrained=false, dir=back, href="#pattern__orchestrator👉 pattern__orchestrator__problem__parallel_programs_", penwidth=2, pos="s,651.36,576.48 663.36,575.52 831.81,561.74 1071,531.6 1271,466 1291.1,459.39 1311.6,451.13 1331.7,442.03", reverse=false, style=tapered, tooltip="*= consistsOf", type=consistsOf]; } }