Šķîþ ţö ḿàîñ çöñţéñţ

É-01: Ţĥé þŕöçéššîñĝ éñĝîñé

Šüḿḿàŕý

Çöñţéñţ îš þŕöçéššéđ ƃý à çĥàññéļ-ƃàšéđ šţŕéàḿîñĝ þîþéļîñé. Éàçĥ ţööļ ŕüñš îñ îţš öŵñ ĝöŕöüţîñé; ţööļš àŕé çöññéçţéđ ƃý ƃüƒƒéŕéđ çĥàññéļš ţĥàţ þŕöṽîđé àüţöḿàţîç ƃàçķþŕéššüŕé. Àñ errgroup.Group çööŕđîñàţéš éŕŕöŕš àñđ þŕöþàĝàţéš çöñţéẋţ çàñçéļļàţîöñ. Ƒöüŕ îñđéþéñđéñţ çöñçüŕŕéñçý ļàýéŕš (îñţŕà-ţööļ ƃļöçķ þàŕàļļéļîšḿ, ƃàţçĥ ƒîļé çöñçüŕŕéñçý, đöçüḿéñţ-ļéṽéļ çöñçüŕŕéñçý, àñđ šţŕéàḿîñĝ öƃšéŕṽàţîöñ) çöḿþöšé ŵîţĥöüţ îñţéŕƒéŕéñçé. Ƒļöŵš àŕé đéçļàŕéđ àš éîţĥéŕ à ĝŕàþĥ öƒ ñöđéš àñđ éđĝéš öŕ à šéǫüéñţîàļ ļîšţ öƒ šţéþš (ŵîţĥ éẋþļîçîţ parallel: ƃļöçķš ƒöŕ ƒàñ-öüţ), ƃöţĥ çöḿþîļéđ ţö ţĥé šàḿé éẋéçüţàƃļé ŕéþŕéšéñţàţîöñ.

Çöñţéẋţ

Ĝö'š ĝöŕöüţîñéš àñđ çĥàññéļš ḿàķé îţ ñàţüŕàļ ţö šţŕüçţüŕé à þîþéļîñé àš çöñçüŕŕéñţ šţàĝéš çöññéçţéđ ƃý ţýþéđ çĥàññéļš. Šüçĥ à þîþéļîñé ḿîẋéš ÇÞÜ-ƃöüñđ ŵöŕķ (ƒöŕḿàţ þàŕšîñĝ, çĥéçķš) ŵîţĥ ÎÖ-ƃöüñđ ŵöŕķ (ĻĻḾ çàļļš, ḿéḿöŕý ļööķüþš). Ţĥé šàḿé þîþéļîñé ḿüšţ àļšö ƃé đŕîṽéñ àţ šéṽéŕàļ šçàļéš: à šîñĝļé ƒîļé öñ à ļàþţöþ; ĥüñđŕéđš öƒ ƒîļéš îñ à ƃàţçĥ; à ļöñĝ-ļîṽéđ þŕöĵéçţ ŵîţĥ ḿàñý đöçüḿéñţš þŕöçéššéđ îñ þàŕàļļéļ.

Îţ ḿüšţ àļšö šüþþöŕţ ƃöţĥ đéçļàŕàţîṽé àüţĥöŕîñĝ (à ṽîšüàļ ƒļöŵ éđîţöŕ, ĥüḿàñ-ŕéàđàƃļé ÝÀḾĻ) àñđ þŕöĝŕàḿḿàţîç çöñšţŕüçţîöñ (ţĥé Ĝö ļîƃŕàŕý) ƒŕöḿ öñé đàţà ḿöđéļ.

Đéçîšîöñ

Çĥàññéļ-ƃàšéđ þîþéļîñé

Çöñţéñţ ƒļöŵš ţĥŕöüĝĥ à çĥàññéļ-ƃàšéđ çöñçüŕŕéñţ þîþéļîñé:

SourcebindingchanTool 1goroutinechanTool 2goroutinechanchanSinkbinding

Çöñţéñţ éñţéŕš ţĥŕöüĝĥ à šöüŕçé ƃîñđîñĝ àñđ ļéàṽéš ţĥŕöüĝĥ à šîñķ ƃîñđîñĝ (É-04). Ƒöŕ ţĥé đéƒàüļţ file ƃîñđîñĝ ţĥéšé àŕé à ĐàţàƑöŕḿàţ ŕéàđéŕ àñđ ŵŕîţéŕ (É-02); à þŕöĵéçţ šţöŕé, à .kpz, öŕ àñ îñţéŕçĥàñĝé ƒîļé ƃîñđ ţĥé šàḿé šţŕéàḿ ŵîţĥ ñö ŕéàđéŕ öŕ ŵŕîţéŕ. Ƃéţŵééñ ţĥé éñđš, éàçĥ ţööļ ŕüñš îñ îţš öŵñ ĝöŕöüţîñé. Ƃüƒƒéŕéđ çĥàññéļš þŕöṽîđé ƃàçķþŕéššüŕé. errgroup.Group çööŕđîñàţéš éŕŕöŕ ĥàñđļîñĝ àçŕöšš ĝöŕöüţîñéš, àñđ çöñţéẋţ çàñçéļļàţîöñ þŕöþàĝàţéš ţö àļļ šţàĝéš.

Þàŕţš çàŕŕý ţýþéđ ŕéšöüŕçéš (Ƒ-02): Ƃļöçķš çöñţàîñ ţŕàñšļàţàƃļé çöñţéñţ, Đàţà çàŕŕîéš šţŕüçţüŕàļ ḿàŕķüþ, Ļàýéŕš ĝŕöüþ ñéšţéđ çöñţéñţ, Ḿéđîà ĥöļđš ƃîñàŕý àššéţš. Ţööļš đéçļàŕé ŵĥîçĥ ŕéšöüŕçé ţýþéš ţĥéý ĥàñđļé; ţĥé ŕéšţ þàšš ţĥŕöüĝĥ üñçĥàñĝéđ.

Éẋéçüţöŕ

flow.DefaultExecutor öŕçĥéšţŕàţéš ţööļ çĥàîñš üšîñĝ ţĥé ĝöŕöüţîñé-þéŕ-ţööļ ḿöđéļ:

  • Éàçĥ ţööļ îñ ţĥé çĥàîñ ŕüñš îñ îţš öŵñ ĝöŕöüţîñé.
  • Ƃüƒƒéŕéđ çĥàññéļš (çöñƒîĝüŕàƃļé šîžé) çöññéçţ àđĵàçéñţ ţööļš.
  • errgroup çöļļéçţš ţĥé ƒîŕšţ éŕŕöŕ àñđ çàñçéļš ţĥé šĥàŕéđ çöñţéẋţ.
  • Þàŕàļļéļ đöçüḿéñţ þŕöçéššîñĝ îš ƃöüñđéđ ƃý à šéḿàþĥöŕé ŵöŕķéŕ þööļ.
  • ToolFactories çŕéàţé ƒŕéšĥ ţööļ îñšţàñçéš þéŕ đöçüḿéñţ, šö çöñçüŕŕéñţ đöçüḿéñţš ñéṽéŕ šĥàŕé ḿüţàƃļé ţööļ šţàţé.

Çöñƒîĝüŕàţîöñ üšéš ţĥé ƒüñçţîöñàļ-öþţîöñš þàţţéŕñ:

executor := flow.NewExecutor(
flow.WithMaxConcurrency(8),
flow.WithChannelSize(128),
flow.WithCollectors(wordCounter, checkReport),
)

flow.WithImmutabilityCheck ţüŕñš öñ ţĥé þéŕ-đîšþàţçĥ çöñţéñţ-îḿḿüţàƃîļîţý ƃàçķšţöþ ƒöŕ à ŕüñ (É-03); îţ îš đéṽ/ţéšţ ţööļîñĝ àñđ îš öƒƒ üñļéšš à çàļļéŕ àšķš ƒöŕ îţ.

Þàŕàļļéļ ƃļöçķ þŕöçéššîñĝ

Ƒöŕ ÎÖ-ƃöüñđ ţööļš (ĻĻḾ ţŕàñšļàţîöñ, ŕéḿöţé çĥéçķš), šéǫüéñţîàļ þéŕ-þàŕţ þŕöçéššîñĝ üñđéŕüţîļîžéš ţĥŕöüĝĥþüţ. tool.ParallelBlockTool ŵŕàþš àñý ţööļ ţö ƒàñ öüţ Ƃļöçķ þŕöçéššîñĝ àçŕöšš Ñ ĝöŕöüţîñéš ŵĥîļé þŕéšéŕṽîñĝ šţŕîçţ Þàŕţ öŕđéŕîñĝ:

InputchanDispatcherseq numberschanfan-out · N goroutines (semaphore-bounded)Worker 1Worker 2Worker NchanReassemblymin-heap · in orderchanOutput

Ţĥé đîšþàţçĥéŕ àššîĝñš ḿöñöţöñîç šéǫüéñçé ñüḿƃéŕš ţö àļļ îñçöḿîñĝ Þàŕţš. Ƃļöçķ Þàŕţš àŕé đîšþàţçĥéđ ţö à šéḿàþĥöŕé-ƃöüñđéđ ŵöŕķéŕ þööļ; ñöñ-Ƃļöçķ Þàŕţš (Đàţà, Ḿéđîà, Ļàýéŕ) þàšš ţĥŕöüĝĥ ţĥé îññéŕ ţööļ šéǫüéñţîàļļý. À ḿîñ-ĥéàþ ŕéàššéḿƃļý ƃüƒƒéŕ çöļļéçţš ŕéšüļţš àñđ éḿîţš ţĥéḿ îñ šţŕîçţ šéǫüéñçé öŕđéŕ, šö đöŵñšţŕéàḿ ţööļš šéé ţĥé šàḿé Þàŕţ öŕđéŕîñĝ ŕéĝàŕđļéšš öƒ ŵĥîçĥ ŵöŕķéŕ ƒîñîšĥéđ ƒîŕšţ.

Àüţö-þàŕàļļéļîšḿ îš à ţööļ þŕöþéŕţý. Éàçĥ ţööļ đéçļàŕéš ToolMeta.DefaultParallelBlocks (É-03); ţĥé ŕüññéŕ ţàķéš ţĥé ḿàẋîḿüḿ àçŕöšš ţĥé ƒļöŵ'š ţööļš àñđ ŵŕàþš éṽéŕý ţööļ àţ ţĥàţ ŵîđţĥ. À þŕöĵéçţ ḿàý þîñ îţš öŵñ ṽàļüé, àñđ --parallel-blocks N öṽéŕŕîđéš ƃöţĥ (--parallel-blocks 1 đîšàƃļéš ţĥé ŵŕàþþéŕ). Ţĥé ţööļš ţĥàţ đéçļàŕé à đéƒàüļţ ţöđàý àŕé ţĥé ĻĻḾ-ƃàçķéđ öñéš, šö àñ öŕđîñàŕý ŕüļéš-öñļý ƒļöŵ ŕüñš šéǫüéñţîàļļý ŵîţĥ ñö çöñƒîĝüŕàţîöñ.

Ƃàţçĥ éẋéçüţöŕ

flow.BatchExecutor þŕöçéššéš ḿüļţîþļé þŕé-ŕéàđ ƒîļéš ţĥŕöüĝĥ à ţööļ çĥàîñ ŵîţĥ çöñƒîĝüŕàƃļé ƒîļé-ļéṽéļ çöñçüŕŕéñçý:

type BatchConfig struct {
FileConcurrency int // max files processed in parallel (default: 1)
ChannelSize int // per-pipeline channel buffer size (default: 64)
SharedResources []io.Closer // resources shared across files (closed at end)
FailFast bool // cancel remaining on first error (default: true)
}

Éàçĥ ƒîļé ĝéţš ƒŕéšĥ ţööļ îñšţàñçéš ƒŕöḿ ŢööļƑàçţöŕý ƒüñçţîöñš, þŕéṽéñţîñĝ šţàţé ļéàķàĝé ƃéţŵééñ çöñçüŕŕéñţ đöçüḿéñţš. Ŕéšüļţš àŕé ŕéţüŕñéđ îñ îñþüţ ƒîļé öŕđéŕ ŕéĝàŕđļéšš öƒ çöḿþļéţîöñ öŕđéŕ. Çöļļéçţöŕš àŕé çàļļéđ ŵîţĥ ḿüţéẋ þŕöţéçţîöñ ƒöŕ ţĥŕéàđ-šàƒé àĝĝŕéĝàţîöñ àçŕöšš ƒîļéš.

Çöñçüŕŕéñçý ļàýéŕîñĝ

Ƒöüŕ îñđéþéñđéñţ çöñçüŕŕéñçý ļàýéŕš çöḿþöšé ŵîţĥöüţ îñţéŕƒéŕéñçé:

ĻàýéŕŠçöþéÇöñţŕöļÖŕđéŕ
ÞàŕàļļéļƂļöçķŢööļƂļöçķš ŵîţĥîñ öñé ţööļÑ ĝöŕöüţîñéš þéŕ ţööļŠţŕîçţ Þàŕţ öŕđéŕ
ƂàţçĥÉẋéçüţöŕḾüļţîþļé ƒîļéšƑîļéÇöñçüŕŕéñçý šéḿàþĥöŕéƑîļé öŕđéŕ þŕéšéŕṽéđ
ÉẋéçüţöŕḾüļţîþļé đöçüḿéñţšḾàẋÇöñçüŕŕéñçý šéḿàþĥöŕéĐöçüḿéñţ öŕđéŕ þŕéšéŕṽéđ
ŢàþþîñĝŢööļÖƃšéŕṽàţîöñÎñļîñé (ñö éẋţŕà ĝöŕöüţîñé)Šéǫüéñţîàļ

Çöļļéçţöŕš àñđ šţŕéàḿîñĝ çöļļéçţöŕš

Çöļļéçţöŕš àĝĝŕéĝàţé ŕéšüļţš àçŕöšš đöçüḿéñţš (ŵöŕđ çöüñţš, çĥéçķ ŕéþöŕţš, ţéŕḿ ļîšţš). Ţĥéý îḿþļéḿéñţ Collect(ctx, item, parts), çàļļéđ àƒţéŕ éàçĥ đöçüḿéñţ çöḿþļéţéš, àñđ Result() ƒöŕ ţĥé ƒîñàļ àĝĝŕéĝàţé. Çöļļéçţöŕš ḿüšţ ƃé ţĥŕéàđ-šàƒé šîñçé ḿüļţîþļé đöçüḿéñţš ḿàý çöḿþļéţé çöñçüŕŕéñţļý.

StreamingCollector éẋţéñđš Collector ŵîţĥ Observe(part) ƒöŕ îñļîñé öƃšéŕṽàţîöñ ŵîţĥöüţ àđđîñĝ à þîþéļîñé šţàĝé. TappingTool ŵŕàþš à ţööļ àñđ îţš šţŕéàḿîñĝ çöļļéçţöŕ: öüţþüţ Þàŕţš àŕé îñţéŕçéþţéđ àñđ þàššéđ ţö Observe() šýñçĥŕöñöüšļý ƃéƒöŕé ƒöŕŵàŕđîñĝ đöŵñšţŕéàḿ. Ţĥîš éñàƃļéš ŕéàļ-ţîḿé ḿéţŕîçš ŵîţĥöüţ ƃüƒƒéŕîñĝ ţĥé éñţîŕé ŕéšüļţ šéţ.

Ƒļöŵ ţŕàçîñĝ àñđ ṽîšüàļîžàţîöñ

flow.TraceRecorder çàþţüŕéš ţîḿéšţàḿþéđ éṽéñţš đüŕîñĝ ƒļöŵ éẋéçüţîöñ. flow.TracingTool ŵŕàþš éàçĥ ţööļ îñ ţĥé çĥàîñ àñđ ŕéçöŕđš éñţéŕ/éẋîţ éṽéñţš ŵîţĥ Þàŕţ šñàþšĥöţš. Ţĥé --trace path/to/trace.json ƒļàĝ öñ kapi run éñàƃļéš ţŕàçîñĝ. Ţĥé öüţþüţ îš à FlowTrace ĴŠÖÑ ƒîļé çöñţàîñîñĝ:

  • Ñöđéš: ţĥé ţööļ çĥàîñ ŵîţĥ çöñçüŕŕéñçý ḿéţàđàţà
  • Éṽéñţš: ţîḿéšţàḿþéđ éñţéŕ/éẋîţ éṽéñţš þéŕ Þàŕţ
  • Þàŕţ šñàþšĥöţš: Þàŕţ šţàţé ƃéƒöŕé àñđ àƒţéŕ éàçĥ ñöđé
  • Đüŕàţîöñ: ţöţàļ ƒļöŵ éẋéçüţîöñ ţîḿé îñ ḿîçŕöšéçöñđš

À ƃŕöŵšéŕ-ƃàšéđ ṽîšüàļîžàţîöñ ŕéñđéŕš ţĥé ţŕàçé àš àñ àñîḿàţéđ ƒļöŵ đîàĝŕàḿ ŵîţĥ þàŕţîçļéš ḿöṽîñĝ ţĥŕöüĝĥ ñöđéš, çĥàññéļ ƒîļļ îñđîçàţöŕš, àñđ ŵöŕķéŕ ļàñé šéþàŕàţîöñ ƒöŕ þàŕàļļéļ ţööļš. Ţĥé þļàýƃàçķ éñĝîñé šüþþöŕţš ṽàŕîàƃļé-šþééđ ŕéþļàý àñđ šééķîñĝ.

Öƃšéŕṽàţîöñ šéàḿ

Ţŕàçîñĝ ţö à ƒîļé îš öñé çöñšüḿéŕ öƒ à ḿöŕé ĝéñéŕàļ šéàḿ. core/observe đéƒîñéš à Tracer / Span îñţéŕƒàçé ŵîţĥ à ñö-öþ đéƒàüļţ àñđ ñö đéþéñđéñçý öñ àñý ţéļéḿéţŕý ļîƃŕàŕý, šö îţ çöšţš ñöţĥîñĝ üñţîļ à ĥöšţ ŕéĝîšţéŕš à ţŕàçéŕ àţ šţàŕţüþ; kapi àñđ ţĥé đéšķţöþ ŕéĝîšţéŕ ñöñé. flow.WrapWithSpans ŵŕàþš éàçĥ ţööļ îñ à çĥàîñ šö ţĥàţ öñé šþàñ öþéñš þéŕ ţööļ îñṽöçàţîöñ (flow.tool), àñđ ţĥé ƒîļé ŕüññéŕ öþéñš öñé þéŕ ƒöŕḿàţ ŕéàđ àñđ ŵŕîţé (format.read, format.write). À SessionTool ķééþš îţš šéššîöñ þàţĥ ţĥŕöüĝĥ ţĥé ŵŕàþþéŕ, šö à ŕéšüḿàƃļé ŕüñ šţîļļ ŕéšüḿéš. Ḿöđéļ þŕöṽîđéŕ çàļļš çàŕŕý ţĥéîŕ öŵñ šþàñ ţĥŕöüĝĥ ţĥé šàḿé šéàḿ (É-07). Öñé šþàñ þéŕ ţööļ ŕàţĥéŕ ţĥàñ þéŕ Þàŕţ ķééþš ţĥé çöšţ þŕöþöŕţîöñàļ ţö ţĥé ĥàñđƒüļ öƒ ţööļš à ƒļöŵ ĥàš, ñöţ ţĥé ţĥöüšàñđš öƒ Þàŕţš îţ ḿöṽéš.

Ƒļöŵ đéƒîñîţîöñš

flow.FlowDefinition îš à ĴŠÖÑ/ÝÀḾĻ-šéŕîàļîžàƃļé šţŕüçţ ţĥàţ çàþţüŕéš à ƒļöŵ ĝŕàþĥ (ñöđéš + éđĝéš) àñđ ţĥé ţööļ çöñƒîĝüŕàţîöñš ñééđéđ ţö ŕéçöñšţŕüçţ à ŕüññàƃļé ƒļöŵ. Ţĥîš šéþàŕàţéš ţĥé đéçļàŕàţîṽé đéšçŕîþţîöñ öƒ à ƒļöŵ ƒŕöḿ îţš ŕüñţîḿé éẋéçüţîöñ.

Éàçĥ FlowNode ĥàš:

  • ÎĐ: üñîǫüé îđéñţîƒîéŕ ŵîţĥîñ ţĥé đéƒîñîţîöñ
  • Ţýþé: flow.NodeTool ƒöŕ à þŕöçéššîñĝ šţéþ. NodeReader àñđ NodeWriter ŕéḿàîñ îñ ţĥé NodeType ṽöçàƃüļàŕý ƒöŕ ţĥé éđîţöŕ'š ĝŕàþĥ, ƃüţ éṽéŕý ƃüîļţ-îñ àñđ šţéþš-çöḿþîļéđ ƒļöŵ çàŕŕîéš ţööļ ñöđéš öñļý, ƃéçàüšé ţĥé éñđš àŕé ƃîñđîñĝš (É-04)
  • Ñàḿé: ţĥé ŕéĝîšţéŕéđ ñàḿé öƒ ţĥé ţööļ (é.ĝ. "pseudo-translate")
  • Ļàƃéļ: öþţîöñàļ đîšþļàý ļàƃéļ ƒöŕ ÜÎ ŕéñđéŕîñĝ
  • Çöñƒîĝ: öþţîöñàļ ķéý-ṽàļüé çöñƒîĝüŕàţîöñ ḿàþ
  • Þöšîţîöñ: ẋ/ý çööŕđîñàţéš ƒöŕ ṽîšüàļ ļàýöüţ îñ ţĥé ƒļöŵ éđîţöŕ

Ƃîñđîñĝš (É-04). À ƒļöŵ'š šöüŕçé àñđ šîñķ àŕé ƃîñđîñĝš ŕéšöļṽéđ ƒŕöḿ îñṽöçàţîöñ çöñţéẋţ (ƒîļé, ţĥé þŕöĵéçţ šţöŕé, à .kpz, îñţéŕçĥàñĝé îḿþöŕţ/éẋþöŕţ, öŕ ñöñé), šö ţĥé šàḿé ƒļöŵ ŕüñš öṽéŕ àñý öŕîĝîñ. Éṽéŕý ƃüîļţ-îñ ƒļöŵ ĝŕàþĥ çàŕŕîéš ţööļ ñöđéš öñļý. Ţĥé ĝŕàþĥ îš çöḿþöšîţîöñ; à šîñĝļé ţööļ îš îñṽöķéđ đîŕéçţļý, ñöţ ŵŕàþþéđ îñ à öñé-ţööļ ƒļöŵ.

Éàçĥ FlowEdge çöññéçţš à šöüŕçé ñöđé ţö à ţàŕĝéţ ñöđé. TopologicalOrder() çöḿþüţéš ţĥé éẋéçüţîöñ öŕđéŕ üšîñĝ Ķàĥñ'š àļĝöŕîţĥḿ, ŕéţüŕñîñĝ àñ éŕŕöŕ îƒ à çýçļé îš đéţéçţéđ, šö îñṽàļîđ ƒļöŵ ĝŕàþĥš ñéṽéŕ ŕéàçĥ ţĥé ŕüñţîḿé éẋéçüţöŕ.

Ţĥé ƃüîļţ-îñ ƒļöŵ çàţàļöĝ îš à þŕöđüçţ çöñçéŕñ àñđ ļîṽéš îñ ţĥé ĥöšţ ḿöđüļé (host/flowdef.BuiltInFlows), ñöţ îñ ţĥé éñĝîñé. Îţ çöṽéŕš ţŕàñšļàţîöñ ŵîţĥ ĝüàŕđŕàîļš, çĥéçķš, ḿéḿöŕý ŕéüšé, þšéüđö-ţŕàñšļàţîöñ, ţĥé ḿéđîà-ţö-šüƃţîţļé þàţĥš, àñđ ţĥé ŕéđàçţîöñ-ƃŕàçķéţéđ ƒļöŵš; ţĥé àüţĥöŕîţàţîṽé ļîšţ ŵîţĥ éàçĥ ƒļöŵ'š šţéþš îš kapi flows àñđ ţĥé ĝéñéŕàţéđ çöḿḿàñđ ŕéƒéŕéñçé.

kapi flows ļîšţš öñļý ţĥé çöḿþöšéđ (ḿüļţî-ţööļ) ƃüîļţ-îñ ƒļöŵš, ƃéçàüšé šîñĝļé-ţööļ đéƒîñîţîöñš àŕé šüŕƒàçéđ àš ţöþ-ļéṽéļ ţööļ çöḿḿàñđš ŕàţĥéŕ ţĥàñ àš ƒļöŵš. host/flowdef.FlowStore þéŕšîšţš üšéŕ-çŕéàţéđ ƒļöŵ đéƒîñîţîöñš àš ĴŠÖÑ ƒîļéš öñ đîšķ, đîšţîñĝüîšĥéđ ƃý šöüŕçé:

  • built-in: šĥîþš ŵîţĥ neokapi, îḿḿüţàƃļé
  • user: çŕéàţéđ ƃý ţĥé üšéŕ, šţöŕéđ îñ ţĥé üšéŕ'š çöñƒîĝ đîŕéçţöŕý

À þŕöĵéçţ'š öŵñ ñàḿéđ ƒļöŵš àŕé đéçļàŕéđ îñ îţš ŕéçîþé'š flows: ƃļöçķ ŕàţĥéŕ ţĥàñ šţöŕéđ àš ƒîļéš (É-04).

Šţéþš-ƃàšéđ ÝÀḾĻ ƒöŕḿàţ

À ĥüḿàñ-ƒŕîéñđļý šţéþš ƒöŕḿàţ îš ţĥé þŕîḿàŕý àüţĥöŕîñĝ šüŕƒàçé ƒöŕ ƒļöŵš îñ ÝÀḾĻ (É-03):

apiVersion: v1
kind: FlowDefinition
metadata:
name: Production Pipeline
spec:
steps:
- tool: recycle
config: { fuzzyThreshold: 75 }
- tool: translate
config: { provider: anthropic }
- tool: qa

Šţéþš àŕé šéǫüéñţîàļ ƃý đéƒàüļţ. parallel: ƃļöçķš þŕöṽîđé ƒàñ-öüţ. Ţĥé þàŕšéŕ àüţö-đéţéçţš ţĥé šĥàþé (šţéþš ṽš ĝŕàþĥ) àñđ çöḿþîļéš šţéþš ţö ñöđéš àñđ éđĝéš. Ƃöţĥ þŕöđüçé ţĥé šàḿé ŕüññàƃļé éẋéçüţöŕ.

Ţĥé šţéþš çàŕŕý öñļý ţĥé çöḿþöšîţîöñ. À ƒļöŵ'š šöüŕçé àñđ šîñķ àŕé ƃîñđîñĝš ŕéšöļṽéđ àţ îñṽöçàţîöñ (ƒîļé, ţĥé þŕöĵéçţ šţöŕé, à .kpz, îñţéŕçĥàñĝé, öŕ ñöñé; É-04) ŕàţĥéŕ ţĥàñ ƒîéļđš öƒ ţĥé ƒļöŵ đöçüḿéñţ.

Ƒàñ-öüţ àñđ ƃàţçĥîñĝ

tool.Tee() çöþîéš þàŕţš ţö Ñ öüţþüţ çĥàññéļš, éñàƃļîñĝ ƒàñ-öüţ ţöþöļöĝîéš ŵĥéŕé öñé ñöđé ƒééđš ḿüļţîþļé þàŕàļļéļ ƃŕàñçĥéš. Ţĥé batch ţööļ çöļļéçţš ƃļöçķš îñţö çöñƒîĝüŕàƃļé ƃàţçĥéš ƃéƒöŕé ƒöŕŵàŕđîñĝ, ŵĥîçĥ šüîţš ƃàţçĥ-çàþàƃļé ŕéḿöţé ÀÞΚ àñđ ĻĻḾ þŕöḿþţš ţĥàţ ƃéñéƒîţ ƒŕöḿ ḿüļţîþļé îñþüţš þéŕ ŕéǫüéšţ (Ḿ-05).

Šçŕîþţ šţéþ

Ţĥé script ţööļ ŕüñš üšéŕ-þŕöṽîđéđ ĴàṽàŠçŕîþţ (ÉŠ5) ṽîà ţĥé ĝöĵà ŕüñţîḿé. Éàçĥ ţööļ îñšţàñçé öŵñš îţš öŵñ goja.Runtime, ŵĥîçĥ îš šàƒé ƃéçàüšé ToolFactory ĝîṽéš öñé îñšţàñçé þéŕ ĝöŕöüţîñé. Ţĥé ĴŠ ÀÞÎ éẋþöšéš part, emit(), skip(), àñđ log() ƒöŕ ƒîļţéŕîñĝ àñđ ţŕàñšƒöŕḿîñĝ þàŕţš: ļîĝĥţŵéîĝĥţ çüšţöḿ ţŕàñšƒöŕḿàţîöñš ŵîţĥöüţ Ĝö çöđé. script îš îñ ţĥé éẋéç çļàšš: à ŕéçîþé çàññöţ àŕḿ îţ šîļéñţļý (É-06).

Ţéŕḿîñöļöĝý: Okapi → neokapi

Ƒöŕ ŕéàđéŕš ƒàḿîļîàŕ ŵîţĥ ţĥé Okapi Ƒŕàḿéŵöŕķ, ţĥé éñĝîñé ḿàþš ţö Okapi çöñçéþţš àš ƒöļļöŵš:

Okapi (Ĵàṽà)neokapi (Ĝö)
ƑîļţéŕĐàţàƑöŕḿàţ (Ŕéàđéŕ/Ŵŕîţéŕ)
ŠţéþŢööļ
ÞîþéļîñéƑļöŵ
ÞîþéļîñéĐŕîṽéŕÉẋéçüţöŕ
ÉṽéñţÞàŕţ
ŢéẋţÜñîţƂļöçķ
ŢéẋţƑŕàĝḿéñţŔüñ šéǫüéñçé ([]Run)
ÇöđéŔüñ
ŠţàŕţŠüƃĐöçüḿéñţ/ŠţàŕţŠüƃƑîļţéŕÇĥîļđ Ļàýéŕ

Çöñšéǫüéñçéš

  • Éàçĥ ţööļ ŕüñš çöñçüŕŕéñţļý; ḿüļţî-çöŕé ÇÞÜš àŕé üţîļîžéđ ŵîţĥîñ à šîñĝļé đöçüḿéñţ'š þîþéļîñé.
  • Ḿüļţîþļé đöçüḿéñţš þŕöçéšš îñ þàŕàļļéļ, ƃöüñđéđ ƃý MaxConcurrency.
  • Ƃàçķþŕéššüŕé îš àüţöḿàţîç: à šļöŵ ţööļ çàüšéš îţš îñþüţ çĥàññéļ ţö ƒîļļ, ŵĥîçĥ ƃļöçķš ţĥé üþšţŕéàḿ ţööļ ŵîţĥöüţ ḿàñüàļ çööŕđîñàţîöñ.
  • Çöñţéẋţ çàñçéļļàţîöñ çļéàñļý þŕöþàĝàţéš ţĥŕöüĝĥ ţĥé éñţîŕé çĥàîñ.
  • ŢööļƑàçţöŕîéš éñšüŕé ñö šĥàŕéđ ḿüţàƃļé šţàţé ƃéţŵééñ þàŕàļļéļ đöçüḿéñţš.
  • Çöļļéçţöŕš þŕöṽîđé çŕöšš-đöçüḿéñţ àĝĝŕéĝàţîöñ ŵîţĥöüţ ƃŕéàķîñĝ ţĥé šţŕéàḿîñĝ ḿöđéļ.
  • Ţööļ àüţĥöŕš đö ñöţ ḿàñàĝé ĝöŕöüţîñéš; ţĥé éẋéçüţöŕ ĥàñđļéš ļîƒéçýçļé, àñđ ParallelBlockTool šüþþļîéš îñţŕà-ţööļ þàŕàļļéļîšḿ ŵîţĥöüţ àñý çöñçüŕŕéñçý çöđé îñ ţĥé ţööļ.
  • StreamingCollector éñàƃļéš ŕéàļ-ţîḿé öƃšéŕṽàţîöñ öƒ þîþéļîñé öüţþüţ ŵîţĥöüţ ḿöđîƒýîñĝ ţĥé Þàŕţ šţŕéàḿ öŕ àđđîñĝ ƃüƒƒéŕîñĝ šţàĝéš.
  • Ƒļöŵ ţŕàçîñĝ éñàƃļéš þöšţ-ĥöç đéƃüĝĝîñĝ àñđ ṽîšüàļîžàţîöñ, ĥéļþîñĝ üšéŕš üñđéŕšţàñđ ţööļ ƃéĥàṽîöüŕ àñđ îđéñţîƒý ƃöţţļéñéçķš.
  • TopologicalOrder ṽàļîđàţîöñ çàţçĥéš çýçļéš ƃéƒöŕé ŕüñţîḿé, ĝîṽîñĝ ƒàšţ ƒééđƃàçķ đüŕîñĝ ƒļöŵ àüţĥöŕîñĝ.
  • ĴŠÖÑ àñđ ÝÀḾĻ šéŕîàļîžàţîöñ šüþþöŕţš îḿþöŕţ/éẋþöŕţ àñđ ṽéŕšîöñ çöñţŕöļ öƒ ƒļöŵ çöñƒîĝüŕàţîöñš.
  • Šţéþš-ƃàšéđ ÝÀḾĻ ḿàķéš ƒļöŵ àüţĥöŕîñĝ àççéššîƃļé ŵîţĥöüţ Ĝö; ţĥé ṽîšüàļ éđîţöŕ àñđ ţĥé ÝÀḾĻ šţàý îñ šýñç ƃéçàüšé ƃöţĥ çöḿþîļé ţö ţĥé šàḿé ĝŕàþĥ.

Ŕéļàţéđ