id: aml-transaction-monitoring-triage
namespace: company.fintech
description: |
Automated Anti-Money Laundering (AML) transaction monitoring pipeline.
Ingests daily transactions, calculates velocity scoring using DuckDB,
and routes high-risk structuring alerts to compliance teams.
triggers:
- id: scheduled_aml_sweep
type: io.kestra.plugin.core.trigger.Schedule
description: Nightly schedule to sweep previous day transactions for velocity
and structuring attacks. Shipped disabled by default.
cron: "0 1 * * *"
disabled: true
- id: event_webhook
type: io.kestra.plugin.core.trigger.Webhook
description: Authenticated webhook for batch settlement ingestion or streaming triggers.
key: "{{ secret('WEBHOOK_KEY') }}"
inputs:
- id: severity_threshold
type: FLOAT
defaults: 0.8
description: The threshold score above which a transaction is flagged for
immediate triage.
tasks:
- id: fetch_transaction_batch
type: io.kestra.plugin.scripts.python.Script
description: Simulates recent transaction batches and injects structured
smurfing patterns.
taskRunner:
type: io.kestra.plugin.core.runner.Process
script: |
import csv
import random
from datetime import datetime, timedelta
# Simulate recent transactions (user_id, amount, timestamp, merchant_category, country)
users = [f"U{1000+i}" for i in range(50)]
categories = ["retail", "crypto", "gambling", "transfer", "atm"]
countries = ["US", "UK", "CY", "MT", "KY"]
now = datetime.utcnow()
transactions = []
# Generate normal background volume
for _ in range(500):
u = random.choice(users)
amt = round(random.uniform(5.0, 500.0), 2)
dt = now - timedelta(hours=random.uniform(0, 24))
transactions.append((u, amt, dt.isoformat(), random.choice(["retail", "transfer"]), "US"))
# Inject an AML anomaly (Smurfing / Velocity Attack)
suspect = "U1042"
for _ in range(15):
amt = round(random.uniform(9000.0, 9999.0), 2) # Just under 10k reporting limit
dt = now - timedelta(minutes=random.uniform(0, 60))
transactions.append((suspect, amt, dt.isoformat(), "crypto", random.choice(["CY", "KY", "MT"])))
with open("transactions.csv", "w", newline="") as f:
writer = csv.writer(f)
writer.writerow(["user_id", "amount", "timestamp", "category", "country"])
writer.writerows(transactions)
print("Generated transactions.csv with mock data")
outputFiles:
- transactions.csv
- id: calculate_velocity_risk
type: io.kestra.plugin.jdbc.duckdb.Query
description: Executes in-memory DuckDB queries to compute account transaction
velocity, volume, and risk scores.
communityExtensions: []
inputFiles:
transactions.csv: "{{ outputs.fetch_transaction_batch.outputFiles['transactions.csv'] }}"
sql: |
WITH raw_tx AS (
SELECT * FROM read_csv_auto('transactions.csv', header = true)
),
user_metrics AS (
SELECT
user_id,
COUNT(*) as tx_count,
SUM(amount) as total_volume,
COUNT(DISTINCT country) as unique_countries,
SUM(CASE WHEN category IN ('crypto', 'gambling') THEN amount ELSE 0 END) as high_risk_volume,
MAX(amount) as max_tx_amount
FROM raw_tx
GROUP BY user_id
),
risk_scores AS (
SELECT
user_id,
tx_count,
total_volume,
unique_countries,
high_risk_volume,
CASE
WHEN tx_count > 10 AND high_risk_volume > 50000 AND unique_countries > 1 THEN 1.0
WHEN tx_count > 5 AND total_volume > 20000 THEN 0.6
ELSE 0.1
END as aml_risk_score
FROM user_metrics
)
SELECT
user_id,
tx_count,
total_volume,
unique_countries,
high_risk_volume,
aml_risk_score
FROM risk_scores
WHERE aml_risk_score >= {{ inputs.severity_threshold }}
ORDER BY aml_risk_score DESC, total_volume DESC;
fetchType: FETCH
- id: triage_gate
type: io.kestra.plugin.core.flow.If
description: Evaluates whether any accounts breached the AML severity threshold
and routes to compliance triage.
condition: "{{ outputs.calculate_velocity_risk.rows | length > 0 }}"
then:
- id: alert_compliance
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alerts compliance on-call channel if high-risk accounts are detected.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"channel": "#compliance-alerts",
"text": "🚨 *AML Velocity Alert* 🚨\n\nIdentified `{{ outputs.calculate_velocity_risk.rows | length }}` account(s) exhibiting suspicious structuring or velocity.\n\n*Top Suspect:* `{{ outputs.calculate_velocity_risk.rows[0].user_id }}` (Score: `{{ outputs.calculate_velocity_risk.rows[0].aml_risk_score }}`).\n\nInvestigate immediately in the AML Case Manager."
}
else:
- id: log_clean_batch
type: io.kestra.plugin.core.log.Log
description: Records clean audit log when no transactions breach risk thresholds.
message: "Batch processed successfully. No high-risk transactions detected."
errors:
- id: alert_on_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alerts data engineering team if the triage pipeline encounters an
unhandled error.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"channel": "#data-engineering",
"text": "❌ *AML Triage Pipeline Failed* ❌\nWorkflow `{{ flow.id }}` failed on execution `{{ execution.id }}`."
}