Skip to main content

Creating conditional flows in TPL Dataflow with LinkTo predicates

While building a data processing pipeline with TPL Dataflow, I needed to route messages to different blocks based on specific conditions. The LinkTo method's predicate parameter is the feature I needed to create branching logic in my dataflow network.

In this post, I explore how to use predicates to build conditional flows that are both efficient and maintainable.

Understanding LinkTo predicates

The LinkTo method in TPL Dataflow connects a source block to a target block, creating a pipeline for data to flow through. The method signature includes an optional predicate parameter:

The predicate is a function that evaluates each message and returns true if the message should be sent to the target block, or false if it should be offered to the next linked block in the chain.

A simple example

Let's start with a straightforward example that demonstrates the basic concept. We'll create a pipeline that routes even numbers to one block and odd numbers to another:



How predicate evaluation works

When a message is ready to leave a source block, TPL Dataflow evaluates the predicates in the order the links were created. The first link whose predicate returns true receives the message. If no predicate matches, the message is dropped unless you've set up a catch-all link.

IMPORTANT: This ordering is critical and gives you fine-grained control over message routing.

Creating a ‘catch-all’ path

To ensure no messages are lost, you should always include a default path for messages that don't match any specific condition. You can do this by adding a final link without a predicate:

Real-world example: Processing orders by priority

Here's a more realistic scenario where we process orders based on their priority and value:

Handling unlinked messages

By default, if a message doesn't match any predicate, it's declined and remains in the source block. To handle this gracefully, you have a few options:

So you should always add an alternative path (such as DataflowBlock.NullTarget()) when using a predicate with LinkTo in TPL Dataflow.

Why? When you use a predicate in LinkTo, only items matching the predicate are sent to the target block. Items that do not match are left unprocessed unless you provide another link for them. If you do not handle these items, the source block will never complete, because it waits for all items to be consumed.

More information

Dataflow (Task Parallel Library) - .NET | Microsoft Learn

Popular posts from this blog

Podman– Command execution failed with exit code 125

After updating WSL on one of the developer machines, Podman failed to work. When we took a look through Podman Desktop, we noticed that Podman had stopped running and returned the following error message: Error: Command execution failed with exit code 125 Here are the steps we tried to fix the issue: We started by running podman info to get some extra details on what could be wrong: >podman info OS: windows/amd64 provider: wsl version: 5.3.1 Cannot connect to Podman. Please verify your connection to the Linux system using `podman system connection list`, or try `podman machine init` and `podman machine start` to manage a new Linux VM Error: unable to connect to Podman socket: failed to connect: dial tcp 127.0.0.1:2655: connectex: No connection could be made because the target machine actively refused it. That makes sense as the podman VM was not running. Let’s check the VM: >podman machine list NAME         ...

Cache stampede: when our cache turned against us

While investigating some performance issues, we ran into an ASP.NET Core API that cached a fairly expensive aggregation query for 60 seconds. Under normal load, that was fine: one request rebuilds the cache, everyone else reads from it. Under peak load, dozens of requests would arrive in that same expiry window, all see a cache miss, and all fire the same expensive query in parallel. The database didn't like that. That was the moment when our caching layer stopped helping and started hurting. A burst of requests comes in at the same time, all miss the cache, and all go hammer the database or the downstream API at once. That's a cache stampede . The cache was supposed to protect our backend, and for a few hundred milliseconds it did the opposite. Why this happens IMemoryCache.GetOrCreate (and its async sibling) looks like it protects you, but it doesn't add any locking on its own. Look at the naive version: public async Task<Report> GetReportAsync(string key) ...

A complex system designed from scratch never works

A few years ago, I worked as an architect on a big mainframe rewrite. I still count it as one of my failures. Not because the technology was wrong, but because I couldn't convince the management team to simplify the approach. Years later, the organization is still struggling to get the new system up and running. I left the project at the time, because I couldn't put my name behind an approach that would take very long and cost a lot of money without a working system to show for it along the way. Gall’s Law That memory keeps coming back to me, because it's a textbook case of Gall's Law playing out in real life. Gall's Law , from John Gall's Systemantics , states it plainly: A complex system that works is invariably found to have evolved from a simple system that worked. A complex system designed from scratch never works, and it cannot be patched to make it work. You have to start over with a simple system that works. What does that mean in practice,...