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

> Optimize Flowsint for large-scale graph investigations with Neo4j batch writes, Celery scaling, and Redis caching. Learn techniques for efficient data handling and analysis.

- Repository: [reconurge/flowsint](https://github.com/reconurge/flowsint)
- Tags: performance
- Published: 2026-06-05

---

**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`](https://github.com/reconurge/flowsint/blob/main/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`](https://github.com/reconurge/flowsint/blob/main/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.

```python
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`](https://github.com/reconurge/flowsint/blob/main/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

```python
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`](https://github.com/reconurge/flowsint/blob/main/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.

```cypher
// 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`](https://github.com/reconurge/flowsint/blob/main/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.

```python
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`](https://github.com/reconurge/flowsint/blob/main/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.

```python
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`](https://github.com/reconurge/flowsint/blob/main/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.

```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`](https://github.com/reconurge/flowsint/blob/main/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

```python
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

```python
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`](https://github.com/reconurge/flowsint/blob/main/config.py)** stores Neo4j and Redis connection settings; bump pool size through the driver constructor in **[`vault.py`](https://github.com/reconurge/flowsint/blob/main/vault.py)**.
- **[`utils.py`](https://github.com/reconurge/flowsint/blob/main/utils.py)** provides `chunked`; use it with Cypher `UNWIND` to batch writes into 300–1000 item transactions.
- **[`enums.py`](https://github.com/reconurge/flowsint/blob/main/enums.py)** defines the graph schema; create matching unique constraints and property indexes in Neo4j before bulk ingestion.
- **[`celery.py`](https://github.com/reconurge/flowsint/blob/main/celery.py)** orchestrates enrichment; enable `acks_late=True` and bound retries so failures replay safely.
- Re-use the Redis URL from **[`config.py`](https://github.com/reconurge/flowsint/blob/main/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`](https://github.com/reconurge/flowsint/blob/main/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`](https://github.com/reconurge/flowsint/blob/main/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`](https://github.com/reconurge/flowsint/blob/main/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`](https://github.com/reconurge/flowsint/blob/main/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.