Functional Data Engineering: Replacing Complex SQL with F#
SQL is the standard for data analytics, but its declarative nature has limits. When transformation logic requires complex state management, heterogeneous stream processing, recursive calculations, or non-linear aggregations, SQL queries turn into multi-level CTEs with heavy window functions. This complicates testing, increases compute costs (especially in systems like BigQuery), and makes the code hard to read.
Using functional programming languages, specifically F#, as an intermediate transformation engine (for example, deployed on Google Cloud Run) solves these architectural problems. The F# type system (algebraic data types, discriminated unions) and pattern matching allow you to express complex business logic in compact, strictly typed, and easily testable code.
Below are 10 engineering cases where moving transformations from SQL to the functional paradigm of F# is technically and economically justified.
Case 1: Parsing Heterogeneous Events (Event Sourcing)
Context: Unstructured JSON with dozens of event types (purchases, registrations, page views) enters the data bus (Pub/Sub, Kafka).
SQL Problem: Processing requires a CASE WHEN for each event type and field extraction via JSON_EXTRACT. The resulting table is sparse with many NULL values.
F# Code:
F#
type UserEvent =
| PageView of Url: string * DurationSec: int
| Purchase of OrderId: string * Amount: decimal
| Registration of Email: string
// Transformation function based on Pattern Matching
let processEvent event =
match event with
| PageView (url, duration) when duration > 30 ->
Some { EventType = "HighIntentView"; Value = url }
| PageView _ ->
None
| Purchase (orderId, amount) ->
Some { EventType = "Conversion"; Value = string amount }
| Registration email ->
Some { EventType = "Lead"; Value = email }
Solution Breakdown: Discriminated Unions (DU) are used. The compiler knows all possible variants of UserEvent. The processEvent function applies pattern matching for routing and transformation.
Result: No implicit NULLs. The compiler issues a warning if a new event type is added but not handled in the match expression. Filtering and extraction logic is strictly typed.
Case 2: Stateful Sessionization
Context: Splitting a user clickstream into sessions.
SQL Problem: If the session timeout is dynamic (e.g., 60 minutes for premium users, 30 for regular users), SQL requires complex LAG window functions, time difference calculations, and cumulative sums to generate a session ID.
F# Code:
F#
type SessionState = { CurrentSessionId: string; LastEventTime: DateTime }
let groupIntoSessions (timeoutMinutes: int) (events: (DateTime * string) list) =
let folder state (time, eventData) =
match state with
| None ->
let newId = Guid.NewGuid().ToString()
Some { CurrentSessionId = newId; LastEventTime = time }, [newId, eventData]
| Some s when (time - s.LastEventTime).TotalMinutes > float timeoutMinutes ->
let newId = Guid.NewGuid().ToString()
Some { CurrentSessionId = newId; LastEventTime = time }, [newId, eventData]
| Some s ->
Some { s with LastEventTime = time }, [s.CurrentSessionId, eventData]
events
|> List.sortBy fst
|> List.mapFold folder None
|> snd
Solution Breakdown: The List.mapFold function iterates over the sorted list of events, passing a state accumulator (SessionState). If the time delta exceeds the timeout, a new session ID is generated.
Result: A single-pass algorithm (O(n)). The logic is easily modified for any session break conditions without complicating the code.
Case 3: Data Validation with Error Accumulation (Railway Oriented Programming)
Context: Raw data must undergo complex validation before loading into the DWH.
SQL Problem: When using SQL constraints or WHERE clauses, the row is simply dropped. It is difficult to collect an array of all rejection reasons for a specific row.
F# Code:
F#
type ValidationError = MissingEmail | NegativeAge | InvalidFormat
let validateEmail user =
if String.IsNullOrEmpty user.Email then Error [MissingEmail] else Ok user
let validateAge user =
if user.Age < 0 then Error [NegativeAge] else Ok user
// Applicative validation (accumulating all errors)
let validateUser user =
match validateEmail user, validateAge user with
| Ok _, Ok _ -> Ok user
| Error e1, Error e2 -> Error (e1 @ e2)
| Error e, Ok _ | Ok _, Error e -> Error e
Solution Breakdown: The Result<'T, 'ErrorList> type is used. Instead of exceptions or NULLs, functions return either success or a list of domain errors. The Railway Oriented Programming pattern allows data to pass through a validation pipeline.
Result: Instead of a simple error flag, a comprehensive array of all violated rules for each row is written to the log (or rejected data table).
Case 4: Non-linear Inventory Depletion (FIFO)
Context: Calculating the cost of goods sold using the FIFO (First-In-First-Out) method.
SQL Problem: SQL operates on sets, while FIFO requires strictly sequential, iterative state changes (depleting inventory batches one by one as they are sold). SQL solutions are extremely cumbersome and slow.
F# Code:
F#
type Batch = { BatchId: string; RemainingQty: int; Cost: decimal }
type Sale = { SaleId: string; QtySold: int }
let applyFifo (batches: Batch list) (sale: Sale) =
let rec deplete remainingToSell currentBatches acc =
match currentBatches with
| [] -> acc, [] // Out of stock
| b :: tail when remainingToSell <= b.RemainingQty ->
let updatedBatch = { b with RemainingQty = b.RemainingQty - remainingToSell }
let costRecord = { SaleId = sale.SaleId; BatchId = b.BatchId; Qty = remainingToSell; Cost = b.Cost }
costRecord :: acc, updatedBatch :: tail
| b :: tail ->
let costRecord = { SaleId = sale.SaleId; BatchId = b.BatchId; Qty = b.RemainingQty; Cost = b.Cost }
deplete (remainingToSell - b.RemainingQty) tail (costRecord :: acc)
deplete sale.QtySold batches []
Solution Breakdown: A recursive function iterates through sorted batches, decreasing the balance and collecting cost records until the required sales quantity is covered.
Result: An accurate, fast, and testable algorithm. Transforming a list of sales into a list of cost depletions takes milliseconds.
Case 5: Parsing Semi-structured Logs (Active Patterns)
Context: Extracting metrics from non-standard system text logs.
SQL Problem: Building complex REGEXP_EXTRACT constructs that break with the slightest change in the string format.
F# Code:
F#
open System.Text.RegularExpressions
// Defining an Active Pattern for regex
let (|Regex|_|) pattern input =
let m = Regex.Match(input, pattern)
if m.Success then Some(List.tail [ for g in m.Groups -> g.Value ]) else None
let parseLogLine log =
match log with
| Regex @"ERROR \[(\w+)\] Timeout (\d+)ms" [service; timeout] ->
Some { Level = "ERROR"; Service = service; Metric = int timeout }
| Regex @"INFO \[(\w+)\] Restarted" [service] ->
Some { Level = "INFO"; Service = service; Metric = 0 }
| _ ->
None
Solution Breakdown: Active Patterns allow encapsulating regular expressions into pattern matching elements. The log string is matched directly against the regex, and the captured groups are immediately bound to variables (service, timeout).
Result: Log parsing code becomes declarative. There is no intermediate NULL checking logic—if a regex does not match, the next pattern is evaluated.
Case 6: Complex Conflict Resolution in Deduplication (Merge Rules)
Context: Merging user profiles from different CRM systems into a single Golden Record.
SQL Problem: Endless COALESCE chains, subqueries with source priority logic, and the complexity of handling timestamps for each individual field.
F# Code:
F#
type Field<'a> = { Value: 'a; SourcePriority: int; UpdatedAt: DateTime }
let mergeFields f1 f2 =
if f1.SourcePriority > f2.SourcePriority then f1
elif f1.SourcePriority < f2.SourcePriority then f2
else if f1.UpdatedAt > f2.UpdatedAt then f1 else f2
let mergeProfiles p1 p2 =
{ Email = mergeFields p1.Email p2.Email
Phone = mergeFields p1.Phone p2.Phone
Score = mergeFields p1.Score p2.Score }
Solution Breakdown: Each value is wrapped in a Field<'a> type containing metadata. The mergeFields function implements strict business logic: source priority first, then update date.
Result: The merge logic is isolated and universal for any data type. Adding new rules (e.g., checking string length) is done in a single function.
Case 7: Generating Domain Events from CDC (Change Data Capture)
Context: The input is a stream of database changes (old_state, new_state). The output must be meaningful business events for analytics.
SQL Problem: SQL triggers or views with CASE WHEN old.status != new.status quickly become unmanageable as the number of fields increases.
F# Code:
F#
type OrderState = { Status: string; Price: decimal }
type OrderEvent = PriceDropped of decimal | StatusChanged of string * string
let generateEvents oldState newState =
let events = []
let events =
if oldState.Price > newState.Price then
PriceDropped (oldState.Price - newState.Price) :: events
else events
let events =
if oldState.Status <> newState.Status then
StatusChanged (oldState.Status, newState.Status) :: events
else events
events |> List.rev
Solution Breakdown: The function takes two states (before and after) and returns a list of domain events generated based on the differences. Discriminated unions (OrderEvent) clearly describe what happened.
Result: The transition from a technical database log to an analytical event model occurs at the typed code level.
Case 8: Hierarchical Structures and Trees (BOM)
Context: Unfolding a Bill of Materials (BOM), where assemblies consist of other assemblies and base components.
SQL Problem: Requires WITH RECURSIVE, which consumes a massive amount of memory and is strictly limited by recursion limits in cloud databases.
F# Code:
F#
type Part =
| Component of string * decimal // Name, Price
| Assembly of string * Part list // Name, List of parts
let rec calculateTotalCost part =
match part with
| Component (_, cost) -> cost
| Assembly (_, parts) ->
parts
|> List.map calculateTotalCost
|> List.sum
Solution Breakdown: Algebraic data types are ideal for trees. A Part is either an atomic Component or an Assembly consisting of a list of Part items. A recursive function traverses the tree with tail-call optimization.
Result: A mathematically clean solution with O(n) complexity without loading the database engine.
Case 9: Time Series Interpolation (Gap Filling)
Context: Data from IoT sensors or stock quotes arrives irregularly. A continuous series with a 1-minute step is needed, copying the last known value.
SQL Problem: Requires generating a date array (CROSS JOIN with a calendar) and a LAST_VALUE(...) IGNORE NULLS window function, creating a Cartesian product.
F# Code:
F#
let fillGaps expectedTimes (measurements: (DateTime * float) list) =
let dataMap = Map.ofList measurements
let folder lastKnownValue time =
let currentValue =
match Map.tryFind time dataMap with
| Some v -> v
| None -> lastKnownValue
currentValue, (time, currentValue)
expectedTimes
|> List.mapFold folder 0.0
|> snd
Solution Breakdown: The state (last known value) is passed through List.mapFold along the reference time array. If data for the current minute is missing, the state from the accumulator is used.
Result: A simple linear pass. The need for heavy database join operations is eliminated.
Case 10: Time-Decay Multi-Channel Attribution
Context: Distributing conversion value across a chain of touchpoints, where the weight of a touchpoint depends on how long ago it occurred (exponential decay).
SQL Problem: Embedding mathematical functions like $y = 2^{-t/T}$ inside an SQL query makes it unreadable and difficult for analysts to debug.
F# Code:
F#
type Touchpoint = { Channel: string; DaysBeforeConversion: float }
let calculateTimeDecay halfLife (touchpoints: Touchpoint list) =
let applyDecay tp =
let weight = 2.0 ** (-tp.DaysBeforeConversion / halfLife)
tp.Channel, weight
let weighted = touchpoints |> List.map applyDecay
let totalWeight = weighted |> List.sumBy snd
weighted
|> List.map (fun (ch, w) -> ch, (w / totalWeight)) // Normalize to 100%
Solution Breakdown: Mathematical logic is separated from data storage. The raw weight for each channel is calculated based on the half-life function, and then the weights are normalized.
Result: The algorithm can be covered by unit tests to check edge cases. The code reads like a model specification.
Best Practices for F# Integration
- Separation of Concerns: Do not move basic filtering (
WHERE) or simple grouping (GROUP BY) to F#. Databases (BigQuery, Snowflake) handle this faster and cheaper. The role of F# is processing arrays, structs, state calculation, and non-linear logic. - Deployment Architecture: An optimal pattern for GCP is creating a microservice on Cloud Run. You extract raw aggregates from BigQuery, send the batch to the Cloud Run API (where the F# code runs), and save the transformed JSON/array back to BigQuery or Cloud Storage.
- Unit Testing Transformations: Logic in SQL is tested by creating mock tables and running real queries (dbt tests, Dataform assertions). F# functions are pure functions—they are tested instantly via xUnit/Expecto without provisioning database infrastructure.
- Input and Output Typing: Use F# Type Providers to generate types based on the BigQuery schema or JSON contracts. This ensures the pipeline will fail at compile time if the data schema changes, rather than in production.
Using F# as a data processor allows replacing convoluted, inefficient SQL scripts with strict functional code. Algebraic data types (DU, Records) and Pattern Matching provide 100% coverage of business logic edge cases.
This approach requires a paradigm shift: the data warehouse (DWH) is used exclusively for its primary purpose—storing columnar data and simple aggregations—while heavy math, working with structures, sessions, and graphs is moved to the compute layer based on a functional language. This reduces FinOps costs (by decreasing data scanning in BigQuery) and radically improves the reliability of analytical pipelines.
