Skip to content
Notifications
Clear all

Just built a lead scoring system in Flux, happy to share the flow

2 Posts
2 Users
0 Reactions
1 Views
(@infra_switcher)
Reputable Member
Joined: 4 months ago
Posts: 320
Topic starter   [#29541]

I've seen a dozen posts here asking about building business logic in Flux, usually framed as "can it be done?" The answer is yes, you can hammer a screw with a wrench if you're determined enough. That doesn't mean you should. Flux is a time-series query language, not a general-purpose programming environment. Building a lead scoring system in it is a classic case of using the wrong tool for the job, but since you're going to do it anyway, here's the painful reality and how we duct-taped it together.

The core problem is state. Lead scoring requires tracking entities (leads) over time, applying rules, and updating scores. Flux has no native concept of stateful entities. You're working on streams of data points. Our implementation hinges on `pivot()` to create quasi-rows and `reduce()` to maintain running totals, which is computationally expensive and gets ugly fast.

Here's the skeleton of our scoring rule logic. Each lead interaction is an event written to an `interactions` measurement. We then window and pivot to get a table per lead per scoring period.

```flux
from(bucket: "crm")
|> range(start: -7d)
|> filter(fn: (r) => r["_measurement"] == "lead_interaction")
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|> group(columns: ["lead_id"])
|> reduce(
identity: {score: 0.0, last_email_open: 0, website_visits: 0},
fn: (r, accumulator) => ({
score: accumulator.score + (
// Rule 1: Email open = +5 points
float(v: r.event_type == "email_open" ? 5 : 0) +
// Rule 2: Website visit > 60s = +10 points
float(v: r.event_type == "page_view" and r.duration_seconds > 60 ? 10 : 0) +
// Rule 3: Demo request = +25 points
float(v: r.event_type == "demo_request" ? 25 : 0)
),
last_email_open: if r.event_type == "email_open" then r._time else accumulator.last_email_open,
website_visits: if r.event_type == "page_view" then accumulator.website_visits + 1 else accumulator.website_visits
})
)
|> map(fn: (r) => ({ r with lead_id: r.lead_id, current_score: r.score }))
```

The pain points we immediately encountered:

* **Performance:** Pivoting on high-cardinality lead_id and time series is a resource hog. Our InfluxDB memory usage spiked 40% after deploying this.
* **Debugging:** Tracing why a specific lead has a certain score requires reconstructing the entire reduce sequence. The observability tools for this are non-existent.
* **Rule Complexity:** Adding a simple rule like "deduct 2 points per day of inactivity" required a completely separate join with a `stateDuration` call, which doubled query execution time.
* **Testing:** You cannot unit test a Flux script in isolation. We had to create a full test data pipeline, which defeated the purpose of a "quick" analytics fix.

We built this because the marketing team wanted it "directly on the data" without involving the application backend. It works, but it's a ticking time bomb. The moment they need to join this scored data with, say, CRM opportunity data from PostgreSQL, you're looking at building a custom task with the SQL plugin or exporting via Telegraf, which introduces more moving parts and failure modes.

If you're considering this path, ask yourself if you really want to own a critical business scoring system that lives in your observability database. The maintenance burden and performance tax are significant. This should be a service in your application layer. We're already planning the migration to a proper service for Q3, which will render this entire Flux monstrosity technical debt.

---


Been there, migrated that


   
Quote
(@alexg)
Honorable Member
Joined: 3 months ago
Posts: 564
 

You're absolutely right about the state problem, but the pivot/reduce approach introduces a hidden cost many don't consider: cardinality explosion. Every unique lead becomes a new series, and with windowing, you're materializing that in memory for each scoring period. That's fine for a few hundred leads, but it's a performance cliff waiting to happen. The query engine starts choking not just on computation, but on metadata.

I've seen teams try to scale this pattern and end up with queries that run for minutes because they're battling InfluxDB's series cardinality limits. It feels clever until your monitoring alerts fire because your lead scoring pipeline is consuming more resources than your actual product. There's a reason you don't see this in production systems that handle any real volume.



   
ReplyQuote