You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

TPL Dataflow:如何让ActionBlock在first_source每次发布值时均被调用?

Hey, let's sort out why your ActionBlock only runs once and get it working as expected!

The problem comes down to how JoinBlock behaves: it won't output anything until every one of its input ports has an available element. So if your second_source only sends a single value, it can only pair with the first value from first_source (resulting in 8), and the remaining values (6, 4) get stuck waiting for more input from second_source—that's why your ActionBlock only executes once.

Here are the right approaches depending on your actual needs:


方案1:直接链接first_source到ActionBlock(不需要second_source的情况)

If your goal is just to trigger the ActionBlock every time first_source emits a value, you don't need a JoinBlock at all. Just link the two blocks directly:

TransformBlock<int, int> first_source = new TransformBlock<int, int>(val => val * 2);
ActionBlock<int> actionBlock = new ActionBlock<int>(val => Console.Write($"{val} "));

// Link the blocks and propagate completion so ActionBlock knows when to stop
first_source.LinkTo(actionBlock, new DataflowLinkOptions { PropagateCompletion = true });

// Test input
first_source.Post(4);
first_source.Post(3);
first_source.Post(2);
first_source.Complete();

await actionBlock.Completion;
// Output: 8 6 4

方案2:使用CombineLatestBlock(需要结合second_source's latest value)

If you do need to use values from both sources, and want the ActionBlock to trigger every time first_source has a new value (using the most recent value from second_source), use CombineLatestBlock instead. It emits a new combined result whenever any input source has a new value:

TransformBlock<int, int> first_source = new TransformBlock<int, int>(val => val * 2);
TransformBlock<int, int> second_source = new TransformBlock<int, int>(val => val + 1); // Example logic

// Combine latest values from both sources into a tuple
var combineBlock = new CombineLatestBlock<int, int, (int FirstVal, int SecondVal)>(
    (first, second) => (FirstVal: first, SecondVal: second));

ActionBlock<(int FirstVal, int SecondVal)> actionBlock = new ActionBlock<(int FirstVal, int SecondVal)>(
    tuple => Console.Write($"{tuple.FirstVal} ")); // Output the first value as you wanted

// Link all blocks with completion propagation
first_source.LinkTo(combineBlock, new DataflowLinkOptions { PropagateCompletion = true });
second_source.LinkTo(combineBlock, new DataflowLinkOptions { PropagateCompletion = true });
combineBlock.LinkTo(actionBlock, new DataflowLinkOptions { PropagateCompletion = true });

// Test: Send a value to second_source first, then all values to first_source
second_source.Post(5);
first_source.Post(4);
first_source.Post(3);
first_source.Post(2);

second_source.Complete();
first_source.Complete();

await actionBlock.Completion;
// Output: 8 6 4

方案3:使用BroadcastBlock to reuse second_source's single value

If second_source only ever sends one value, but you want that value to pair with every value from first_source, use a BroadcastBlock to repeat the second_source value for each incoming first_source element:

TransformBlock<int, int> first_source = new TransformBlock<int, int>(val => val * 2);
TransformBlock<int, int> second_source = new TransformBlock<int, int>(val => val + 1);

// BroadcastBlock keeps the latest value and sends it to every linked target
var broadcastSecond = new BroadcastBlock<int>(val => val);
second_source.LinkTo(broadcastSecond, new DataflowLinkOptions { PropagateCompletion = true });

// Now JoinBlock will get a value from broadcastSecond for every first_source element
var joinBlock = new JoinBlock<int, int>();
first_source.LinkTo(joinBlock.Target1, new DataflowLinkOptions { PropagateCompletion = true });
broadcastSecond.LinkTo(joinBlock.Target2, new DataflowLinkOptions { PropagateCompletion = true });

ActionBlock<Tuple<int, int>> actionBlock = new ActionBlock<Tuple<int, int>>(
    tuple => Console.Write($"{tuple.Item1} "));

joinBlock.LinkTo(actionBlock, new DataflowLinkOptions { PropagateCompletion = true });

// Test
second_source.Post(5);
first_source.Post(4);
first_source.Post(3);
first_source.Post(2);

second_source.Complete();
first_source.Complete();

await actionBlock.Completion;
// Output: 8 6 4

内容的提问来源于stack exchange,提问作者Muhammad Ramzan

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 06:41:35