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

