How to Optimize Flowsint for Large-Scale Graph Investigations: Neo4j Batch Writes, Celery Scaling, and Redis Caching

Batch Neo4j writes using the chunked utility, add unique constraints on node identifiers, scale Celery enrichment workers with acks_late=True, and tune connection pools via flowsint-core settings to optimize Flowsint for large-scale graph investigations.

The reconurge/flowsint platform persists investigative metadata in PostgreSQL and graph relationships in Neo4j, orchestrating enrichment tasks through Celery. To optimize Flowsint for large-scale graph investigations, you must minimize per-transaction overhead, parallelize ingestion safely, and cache expensive external lookups. The following sections map each optimization to concrete source files in flowsint-core and flowsint-api.

Tune Neo4j Driver Connection Pooling

Connection defaults are declared in flowsint-core/src/flowsint_core/core/config.py at lines 11–14, while the driver instance is managed inside flowsint-core/src/flowsint_core/core/vault.py. At production scale, instantiate the driver with an expanded pool so concurrent Celery workers do not block waiting for an idle socket.

from neo4j import GraphDatabase
from flowsint_core.core.config import settings

driver = GraphDatabase.driver(
    f"bolt://{settings.NEO4J_HOST}:{settings.NEO4J_PORT}",
    auth=(settings.NEO4J_USER, settings.NEO4J_PASSWORD),
    max_connection_pool_size=50,
    connection_acquisition_timeout=30.0,
)

A pool size of 50 handles dozens of parallel enrichment tasks without exhaustion. Set connection_acquisition_timeout high enough to survive traffic spikes, but monitor driver metrics so requests do not queue indefinitely.

Batch Graph Writes with the chunked Utility

The chunked helper defined in flowsint-core/src/flowsint_core/utils.py at lines 14–21 splits an iterable into fixed-size blocks. Pair it with a Cypher UNWIND statement to merge thousands of nodes in a single transaction.

Bulk Node Insertion Example

from flowsint_core.core.config import settings
from neo4j import GraphDatabase
from flowsint_core.utils import chunked

driver = GraphDatabase.driver(
    f"bolt://{settings.NEO4J_HOST}:{settings.NEO4J_PORT}",
    auth=(settings.NEO4J_USER, settings.NEO4J_PASSWORD),
    max_connection_pool_size=50,
    connection_acquisition_timeout=30.0,
)

def create_nodes(nodes):
    """
    `nodes` is a list of dicts:
    {"id": "...", "type": "domain", "props": {...}}
    """
    cypher = """
    UNWIND $batch AS n
    MERGE (node:%s {iid: n.id})
    SET node += n.props
    """
    with driver.session() as session:
        for batch in chunked(nodes, 500):
            label = batch[0]["type"].upper()
            stmt = cypher % label
            session.run(stmt, batch=batch)

Using batches of 500 keeps each transaction under Neo4j’s default heap limits while eliminating round-trip latency. Increase the batch size only after profiling memory on your target hardware.

Add Unique Constraints and Property Indexes

Node and edge labels are centralized in flowsint-core/src/flowsint_core/core/enums.py. Every investigation node type—DOMAIN, IP, ASN, EMAIL, and others—should carry a unique business key. Create constraints and indexes immediately after deployment so MERGE operations execute index seeks rather than label scans.

// Constraints aligned with NodeType values in enums.py
CREATE CONSTRAINT domain_iid IF NOT EXISTS
    FOR (d:DOMAIN) REQUIRE d.iid IS UNIQUE;

CREATE INDEX domain_name IF NOT EXISTS
    FOR (d:DOMAIN) ON (d.name);

Repeat the pattern for each label defined in NodeType. Unique constraints on iid prevent duplicate entities during concurrent enrichment, while secondary property indexes accelerate filtered reads in the investigation UI.

Scale Parallel Enrichment with Celery Workers

Heavy lookup tasks are dispatched through Celery inside flowsint-core/src/flowsint_core/core/celery.py. Configure workers with acks_late=True and moderate concurrency so a crashed process does not drop graph mutations.

from celery import shared_task
from flowsint_core.utils import chunked

@shared_task(bind=True, acks_late=True, max_retries=3)
def enrich_domain(self, domain_id, enrichment_data):
    try:
        with driver.session() as session:
            for batch in chunked(enrichment_data, 300):
                session.write_transaction(_apply_batch, batch)
    except Exception as exc:
        raise self.retry(exc=exc, countdown=60)

def _apply_batch(tx, batch):
    cypher = """
    UNWIND $batch AS rel
    MATCH (src:DOMAIN {iid: rel.src})
    MATCH (dst:%s {iid: rel.dst})
    MERGE (src)-[r:%s]->(dst)
    SET r += rel.props
    """
    tx.run(cypher, batch=batch)

Processing 300 relationships per transaction balances throughput with lock contention. If you observe Neo4j deadlock exceptions in the worker logs, reduce the batch size or serialize writes by investigation case.

Cache Expensive Lookups Using Redis

flowsint-core/src/flowsint_core/core/config.py exposes the Redis broker URL at line 16 via REDIS_URL. Re-use that same Redis instance to cache external API calls such as WHOIS lookups, preventing redundant network latency and reducing the volume of new nodes created.

import json
import redis
from functools import wraps
from flowsint_core.core.config import settings

redis_client = redis.from_url(settings.REDIS_URL)

def redis_cache(ttl: int = 3600):
    def decorator(fn):
        @wraps(fn)
        def wrapper(*args, **kwargs):
            key = f"{fn.__name__}:{args}:{kwargs}"
            cached = redis_client.get(key)
            if cached:
                return json.loads(cached)
            result = fn(*args, **kwargs)
            redis_client.setex(key, ttl, json.dumps(result))
            return result
        return wrapper
    return decorator

A 24-hour TTL for stable WHOIS records is usually safe. Deduct cached responses from your enrichment pipeline before the Celery task reaches the graph write stage to spare Neo4j unnecessary MERGE work.

Paginate Graph API Endpoints

The frontend investigation UI fetches subgraphs through flowsint-api/app/api/routes/flows.py. Expose skip and limit parameters on every graph endpoint, and push pagination down to Neo4j rather than filtering in Python.

from fastapi import APIRouter

router = APIRouter()

@router.get("/graph")
def get_subgraph(skip: int = 0, limit: int = 100):
    query = """
    MATCH (n)
    RETURN n
    SKIP $skip LIMIT $limit
    """
    with driver.session() as session:
        result = session.run(query, skip=skip, limit=limit)
    return [record["n"] for record in result]

Rendering 100 nodes at a time keeps the browser responsive. Extend this pattern to relationship queries by adding an EdgeType filter that maps directly to the enum values defined in flowsint-core/src/flowsint_core/core/enums.py.

Complete Integration Snippets

The following helpers belong in flowsint-core modules and can be imported by enrichers and API routes alike.

Bulk Merge Helper

from flowsint_core.utils import chunked

def bulk_merge_nodes(driver, nodes):
    """
    `nodes` is a list of dicts:
    {"iid": str, "label": NodeType, "props": dict}
    """
    cypher = """
    UNWIND $batch AS n
    MERGE (node:%s {iid: n.iid})
    SET node += n.props
    """
    with driver.session() as session:
        for batch in chunked(nodes, 1000):
            label = batch[0]["label"].value.upper()
            stmt = cypher % label
            session.run(stmt, batch=batch)

Cached WHOIS Lookup

import json, redis
from flowsint_core.core.config import settings

redis_client = redis.from_url(settings.REDIS_URL)

def cached_whois(domain: str):
    key = f"whois:{domain}"
    cached = redis_client.get(key)
    if cached:
        return json.loads(cached)

    import whois
    data = whois.whois(domain)
    redis_client.setex(key, 86400, json.dumps(data))
    return data

Combine both snippets inside a Celery task that writes only new or changed data back to Neo4j.

Summary

  • config.py stores Neo4j and Redis connection settings; bump pool size through the driver constructor in vault.py.
  • utils.py provides chunked; use it with Cypher UNWIND to batch writes into 300–1000 item transactions.
  • enums.py defines the graph schema; create matching unique constraints and property indexes in Neo4j before bulk ingestion.
  • celery.py orchestrates enrichment; enable acks_late=True and bound retries so failures replay safely.
  • Re-use the Redis URL from config.py line 16 to cache external API calls and avoid redundant node creation.
  • Return paginated responses from flowsint-api/app/api/routes/flows.py so the UI never fetches an entire million-node graph at once.

Frequently Asked Questions

What is the optimal Neo4j batch size when writing Flowsint nodes?

For most hardware profiles, 300 to 500 nodes or edges per UNWIND batch avoids heap pressure and lock contention. If you monitor Neo4j logs and see frequent Neo.TransientError.Transaction.DeadlockDetected messages, lower the chunk size in flowsint-core/src/flowsint_core/utils.py or serialize writes per investigation case.

How does Celery improve large-scale graph investigation performance?

Celery distributes CPU-intensive enrichment tasks across multiple workers. By setting acks_late=True in the task decorator, the broker ensures that a task is only acknowledged after completion, so crashed workers do not lose graph mutations. This keeps Neo4j write queues steady rather than bursty.

Where should Neo4j indexes be created in a Flowsint deployment?

Run CREATE CONSTRAINT and CREATE INDEX statements directly against the Neo4j database after deployment. The constraints should target the iid property for every label defined in flowsint-core/src/flowsint_core/core/enums.py, because the bulk insertion helpers in flowsint-core rely on MERGE against that key.

Can the existing Redis broker cache enrichment results?

Yes. flowsint-core/src/flowsint_core/core/config.py exposes REDIS_URL at line 16 for Celery, and the same URL can initialize a redis client inside enrichment modules. Cache stable external data such as WHOIS records with a TTL of one day to reduce both API quota consumption and Neo4j write load.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →