The Silent Killer of Backtest Integrity

A quant researcher runs a mean-reversion strategy on two years of minute-bar data. The backtest shows a Sharpe ratio of 1.8. Live deployment underperforms by 40%. The culprit? A data pipeline that silently dropped 3% of records during a daily sync — records that, coincidentally, coincided with the strategy's best-performing regime.

This is not a hypothetical. Data synchronization failures are among the most insidious problems in quantitative engineering. They do not throw errors. They do not send alerts. They corrupt your data quietly, and you discover the damage only when your live system diverges from your backtest in ways that cost you real money.

The solution requires a disciplined approach to three problems: avoiding duplicate data on repeated pulls, tracking data source versions to detect backfills, and implementing hash-based integrity checks to catch silent mutations. This article provides production-grade code and architectural patterns for solving all three.


The Three Failure Modes of Scheduled Data Sync

Before diving into solutions, it helps to understand the enemy. Scheduled data synchronization fails in three characteristic ways.

Duplicate insertion occurs when a pipeline runs, fails partway through, retries, and inserts records that already exist. Without a deduplication strategy, your database accumulates identical rows, skewing volume-weighted metrics and inflating apparent trading activity.

Stale data persistence occurs when a data source modifies historical records — a common practice known as backfilling. Corporate actions adjustments, dividend corrections, and exchange data amendments all produce this behavior. A pipeline that blindly inserts only "new" records by timestamp will never learn about these corrections.

Silent mutation is the most dangerous failure mode. The data source changes a record's values without changing its timestamp. Your pipeline sees no "new" data. Your database retains the old values. The discrepancy persists indefinitely.

TickDB provides version metadata and supports high-volume retrieval, but your local pipeline must implement the logic to handle these scenarios. The techniques below apply to any data source; the code examples use TickDB's REST API as the concrete implementation.


Strategy 1: UPSERT — The Foundation of Deduplication

The most reliable deduplication mechanism is the database-level UPSERT — an atomic operation that inserts a record if it does not exist and updates it if it does. Every major database supports this: PostgreSQL's ON CONFLICT, MySQL's INSERT ... ON DUPLICATE KEY UPDATE, and SQLite's INSERT OR REPLACE.

The Unique Key Problem

UPSERT requires a unique constraint to determine "sameness." For market data, the natural unique key is the combination of symbol, timestamp, and granularity. A 5-minute bar for AAPL at 9:30 AM on a specific date is unique by construction.

import os
import requests
import psycopg2
from psycopg2.extras import execute_values
from datetime import datetime, timedelta
import time

# ──────────────────────────────────────────────
# Configuration
# ──────────────────────────────────────────────
TICKDB_API_KEY = os.environ.get("TICKDB_API_KEY")
DB_CONFIG = {
    "host": os.environ.get("DB_HOST", "localhost"),
    "dbname": os.environ.get("DB_NAME", "market_data"),
    "user": os.environ.get("DB_USER", "quant_user"),
    "password": os.environ.get("DB_PASSWORD"),
}

# Table schema for OHLCV bars
CREATE_TABLE_SQL = """
CREATE TABLE IF NOT EXISTS ohlcv_bars (
    symbol      VARCHAR(20) NOT NULL,
    timestamp   TIMESTAMPTZ NOT NULL,
    interval    VARCHAR(10) NOT NULL DEFAULT '5m',
    open        NUMERIC(18, 8),
    high        NUMERIC(18, 8),
    low         NUMERIC(18, 8),
    close       NUMERIC(18, 8),
    volume      BIGINT,
    tick_count  INTEGER,
    created_at  TIMESTAMPTZ DEFAULT NOW(),
    updated_at  TIMESTAMPTZ DEFAULT NOW(),
    PRIMARY KEY (symbol, timestamp, interval)
);
"""

UPSERT_SQL = """
INSERT INTO ohlcv_bars (symbol, timestamp, interval, open, high, low, close, volume, tick_count, updated_at)
VALUES %s
ON CONFLICT (symbol, timestamp, interval)
DO UPDATE SET
    open = EXCLUDED.open,
    high = EXCLUDED.high,
    low = EXCLUDED.low,
    close = EXCLUDED.close,
    volume = EXCLUDED.volume,
    tick_count = EXCLUDED.tick_count,
    updated_at = EXCLUDED.updated_at;
"""


def fetch_ohlcv_batch(symbol: str, start_ts: int, end_ts: int, interval: str = "5m") -> list[dict]:
    """
    Fetch OHLCV bars from TickDB for a given time range.
    start_ts and end_ts are Unix timestamps in milliseconds.
    """
    url = "https://api.tickdb.ai/v1/market/kline"
    headers = {"X-API-Key": TICKDB_API_KEY}
    params = {
        "symbol": symbol,
        "interval": interval,
        "start": start_ts,
        "end": end_ts,
        "limit": 1000,
    }

    response = requests.get(url, headers=headers, params=params, timeout=(3.05, 27))
    response.raise_for_status()
    data = response.json()

    if data.get("code") != 0:
        raise RuntimeError(f"API error {data.get('code')}: {data.get('message')}")

    return data.get("data", [])

Handling Rate Limits in Batched Fetches

When pulling historical data across a large date range, you must paginate and respect rate limits. The code: 3001 response indicates rate limiting; read the Retry-After header and wait before retrying.

def fetch_with_retry(symbol: str, start_ts: int, end_ts: int, interval: str = "5m", max_retries: int = 5):
    """Fetch OHLCV data with exponential backoff and rate-limit handling."""
    url = "https://api.tickdb.ai/v1/market/kline"
    headers = {"X-API-Key": TICKDB_API_KEY}
    params = {
        "symbol": symbol,
        "interval": interval,
        "start": start_ts,
        "end": end_ts,
        "limit": 1000,
    }

    for attempt in range(max_retries):
        response = requests.get(url, headers=headers, params=params, timeout=(3.05, 27))

        if response.status_code == 200:
            data = response.json()
            if data.get("code") == 0:
                return data.get("data", [])
            if data.get("code") == 3001:
                # Rate limited — respect Retry-After
                retry_after = int(response.headers.get("Retry-After", 5))
                print(f"Rate limited. Waiting {retry_after}s before retry {attempt + 1}/{max_retries}")
                time.sleep(retry_after)
                continue
            raise RuntimeError(f"API error {data.get('code')}: {data.get('message')}")

        # HTTP-level retry for 5xx errors
        if response.status_code >= 500 and attempt < max_retries - 1:
            delay = min(2 ** attempt + 0.1 * __import__("random").random(), 30)
            print(f"HTTP {response.status_code}. Retrying in {delay:.1f}s")
            time.sleep(delay)
            continue

        response.raise_for_status()

    raise RuntimeError(f"Failed after {max_retries} attempts")


def sync_symbol_to_db(symbol: str, start_date: datetime, end_date: datetime, interval: str = "5m"):
    """
    Full pipeline: fetch OHLCV from TickDB and UPSERT to local PostgreSQL.
    Processes data in daily batches to handle large ranges efficiently.
    """
    conn = psycopg2.connect(**DB_CONFIG)
    cursor = conn.cursor()

    # Ensure table exists
    cursor.execute(CREATE_TABLE_SQL)
    conn.commit()

    current = start_date
    total_records = 0

    while current < end_date:
        batch_end = min(current + timedelta(days=1), end_date)
        start_ts = int(current.timestamp() * 1000)
        end_ts = int(batch_end.timestamp() * 1000)

        try:
            bars = fetch_with_retry(symbol, start_ts, end_ts, interval)
        except Exception as e:
            print(f"Error fetching {symbol} for {current.date()}: {e}")
            current = batch_end
            continue

        if bars:
            # Prepare records for batch UPSERT
            records = [
                (
                    bar["symbol"],
                    datetime.fromtimestamp(bar["t"] / 1000, tz=__import__("datetime").timezone.utc),
                    interval,
                    bar["o"],
                    bar["h"],
                    bar["l"],
                    bar["c"],
                    bar["v"],
                    bar.get("n", 0),
                    datetime.now(timezone.utc),
                )
                for bar in bars
            ]

            execute_values(cursor, UPSERT_SQL, records)
            conn.commit()
            total_records += len(records)

        current = batch_end

    cursor.close()
    conn.close()
    print(f"Synced {total_records} records for {symbol} from {start_date.date()} to {end_date.date()}")

Why UPSERT Beats INSERT with Deduplication

A naive deduplication approach is to SELECT to check existence before INSERT. This has two problems: it doubles your query volume, and it creates a race condition in concurrent pipelines. UPSERT is atomic and eliminates both issues. The unique constraint acts as your guardrail at the database level — you cannot accidentally insert duplicates regardless of how your application logic behaves.


Strategy 2: Version Number Tracking for Backfill Detection

UPSERT solves duplicate insertion, but it does not solve stale data persistence. If the data source modifies a historical record, your UPSERT will overwrite the old value — but only if you fetch that record again. If your pipeline only fetches "new" data (records after the last sync timestamp), you will never know about corrections to older data.

The solution is to track a version identifier alongside each record. TickDB provides this via the v field in API responses — a version number that increments whenever the record changes.

The Version Tracking Schema

Add a data_version column to your table:

ALTER_TABLE_SQL = """
ALTER TABLE ohlcv_bars ADD COLUMN IF NOT EXISTS data_version BIGINT DEFAULT 0;
"""

UPSERT_WITH_VERSION_SQL = """
INSERT INTO ohlcv_bars (symbol, timestamp, interval, open, high, low, close, volume, tick_count, data_version, updated_at)
VALUES %s
ON CONFLICT (symbol, timestamp, interval)
DO UPDATE SET
    open = EXCLUDED.open,
    high = EXCLUDED.high,
    low = EXCLUDED.low,
    close = EXCLUDED.close,
    volume = EXCLUDED.volume,
    tick_count = EXCLUDED.tick_count,
    data_version = EXCLUDED.data_version,
    updated_at = EXCLUDED.updated_at
WHERE ohlcv_bars.data_version < EXCLUDED.data_version;
"""

The WHERE ohlcv_bars.data_version < EXCLUDED.data_version clause is critical. Without it, UPSERT overwrites regardless of version. With it, the update only occurs if the incoming version is newer. This means:

  • First insert: version 1 → record created.
  • Data source backfills: same timestamp, version 2 → record updated.
  • Re-delivery of same version: version 1 → no change (record already matches).

Implementing a Version-Aware Sync

def get_last_version(conn, symbol: str, interval: str = "5m") -> int:
    """Query the highest data_version we have for a given symbol."""
    with conn.cursor() as cur:
        cur.execute(
            """
            SELECT COALESCE(MAX(data_version), 0)
            FROM ohlcv_bars
            WHERE symbol = %s AND interval = %s
            """,
            (symbol, interval),
        )
        return cur.fetchone()[0]


def sync_with_backfill_detection(symbol: str, start_date: datetime, end_date: datetime, interval: str = "5m"):
    """
    Sync pipeline that detects both new data and backfilled corrections.
    Re-fetches a lookback window to catch version changes.
    """
    conn = psycopg2.connect(**DB_CONFIG)

    # Lookback window: 7 days of history is re-checked on every run
    # to catch any backfilled corrections from the data source.
    lookback_days = 7
    effective_start = max(start_date, datetime.now(timezone.utc) - timedelta(days=lookback_days))

    current = effective_start
    total_updated = 0

    while current < end_date:
        batch_end = min(current + timedelta(days=1), end_date)
        start_ts = int(current.timestamp() * 1000)
        end_ts = int(batch_end.timestamp() * 1000)

        try:
            bars = fetch_with_retry(symbol, start_ts, end_ts, interval)
        except Exception as e:
            print(f"Error fetching {symbol} for {current.date()}: {e}")
            current = batch_end
            continue

        if bars:
            records = [
                (
                    bar["symbol"],
                    datetime.fromtimestamp(bar["t"] / 1000, tz=timezone.utc),
                    interval,
                    bar["o"],
                    bar["h"],
                    bar["l"],
                    bar["c"],
                    bar["v"],
                    bar.get("n", 0),
                    bar["v"],  # data_version from TickDB
                    datetime.now(timezone.utc),
                )
                for bar in bars
            ]

            with conn.cursor() as cur:
                execute_values(cur, UPSERT_WITH_VERSION_SQL, records)
                conn.commit()
                total_updated += len(records)

        current = batch_end

    conn.close()
    print(f"Processed {total_updated} records for {symbol}")

Choosing Your Lookback Window

The lookback window represents a trade-off between data freshness and API call volume. A 7-day lookback is aggressive and catches most corporate action adjustments, which exchanges typically amend within 3–5 business days. For less critical data, a 2-day window reduces API usage. For regulatory or audit data, a 30-day window is appropriate.


Strategy 3: Hash Validation for Silent Mutation Detection

Version numbers solve the backfill problem, but they require the data source to provide versioning. Not all fields change with every version update, and not all data sources expose version metadata. For these cases, hash-based validation provides an additional integrity layer.

The idea is straightforward: compute a hash of the incoming record's meaningful fields. Store the hash alongside the record. If the hash changes on a subsequent fetch but the version number does not, you have detected a silent mutation.

import hashlib
import json


def compute_record_hash(bar: dict) -> str:
    """
    Compute a SHA-256 hash of the canonical fields in an OHLCV bar.
    Excludes metadata fields like timestamps and version numbers.
    """
    canonical = {
        "symbol": bar["symbol"],
        "open": bar["o"],
        "high": bar["h"],
        "low": bar["l"],
        "close": bar["c"],
        "volume": bar["v"],
        "tick_count": bar.get("n", 0),
    }
    serialized = json.dumps(canonical, sort_keys=True)
    return hashlib.sha256(serialized.encode()).hexdigest()


def check_mutation(cursor, symbol: str, timestamp: datetime, interval: str, new_hash: str) -> bool:
    """
    Check if a record's hash has changed since the last sync.
    Returns True if a mutation is detected.
    """
    cursor.execute(
        """
        SELECT data_hash FROM ohlcv_bars
        WHERE symbol = %s AND timestamp = %s AND interval = %s
        """,
        (symbol, timestamp, interval),
    )
    row = cursor.fetchone()

    if row is None:
        return False  # New record — no mutation to detect

    return row[0] != new_hash


def flag_mutations(symbol: str, start_date: datetime, end_date: datetime, interval: str = "5m"):
    """
    Scan a date range and flag any records where the hash has changed.
    Useful for auditing data quality and detecting silent corrections.
    """
    conn = psycopg2.connect(**DB_CONFIG)
    mutations = []

    current = start_date
    while current < end_date:
        batch_end = min(current + timedelta(days=1), end_date)
        start_ts = int(current.timestamp() * 1000)
        end_ts = int(batch_end.timestamp() * 1000)

        try:
            bars = fetch_with_retry(symbol, start_ts, end_ts, interval)
        except Exception as e:
            print(f"Error scanning {symbol} for {current.date()}: {e}")
            current = batch_end
            continue

        with conn.cursor() as cur:
            for bar in bars:
                ts = datetime.fromtimestamp(bar["t"] / 1000, tz=timezone.utc)
                new_hash = compute_record_hash(bar)

                if check_mutation(cur, symbol, ts, interval, new_hash):
                    mutations.append({
                        "symbol": symbol,
                        "timestamp": ts.isoformat(),
                        "interval": interval,
                        "new_hash": new_hash,
                    })

        current = batch_end

    conn.close()

    if mutations:
        print(f"Detected {len(mutations)} mutations for {symbol}:")
        for m in mutations[:10]:  # Print first 10
            print(f"  {m['symbol']} @ {m['timestamp']} [{m['interval']}] — values changed")
    else:
        print(f"No mutations detected for {symbol}")

    return mutations

Hash Validation vs. Version Tracking — When to Use Each

Version tracking is more efficient because it piggybacks on the data source's own change detection. Use version tracking when the data source provides version metadata (as TickDB does). Use hash validation as a supplementary layer when you need to detect changes that version metadata might miss — for example, when merging data from multiple sources where each source has different version semantics.


Putting It Together: A Complete Sync Job

The following is a production-ready daily sync job that ties all three strategies together. It includes scheduling, error handling, logging, and Slack alerting for mutations.

import schedule
import time
import logging
from datetime import datetime, timedelta, timezone
from typing import Optional

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s [%(levelname)s] %(message)s",
    handlers=[
        logging.FileHandler("/var/log/tickdb_sync.log"),
        logging.StreamHandler(),
    ],
)
logger = logging.getLogger(__name__)


class DataSyncJob:
    def __init__(self, symbols: list[str], interval: str = "5m"):
        self.symbols = symbols
        self.interval = interval
        self.webhook_url: Optional[str] = os.environ.get("SLACK_WEBHOOK_URL")

    def run(self):
        """
        Daily sync entry point. Scheduled to run after market close (4:15 PM ET).
        """
        logger.info(f"Starting daily sync for {len(self.symbols)} symbols")
        start_time = datetime.now(timezone.utc)

        total_mutations = []

        for symbol in self.symbols:
            try:
                # Sync last 7 days to catch backfills
                end_date = datetime.now(timezone.utc)
                start_date = end_date - timedelta(days=7)

                sync_with_backfill_detection(symbol, start_date, end_date, self.interval)

                # Audit for mutations on the last 30 days
                audit_end = datetime.now(timezone.utc)
                audit_start = audit_end - timedelta(days=30)
                mutations = flag_mutations(symbol, audit_start, audit_end, self.interval)
                total_mutations.extend(mutations)

            except Exception as e:
                logger.error(f"Failed to sync {symbol}: {e}")

        elapsed = (datetime.now(timezone.utc) - start_time).total_seconds()
        logger.info(f"Daily sync completed in {elapsed:.1f}s. Mutations: {len(total_mutations)}")

        if total_mutations and self.webhook_url:
            self._send_alert(total_mutations)

    def _send_alert(self, mutations: list[dict]):
        """Send a Slack alert when mutations are detected."""
        if not mutations:
            return

        blocks = [
            {
                "type": "header",
                "text": {"type": "plain_text", "text": "⚠️ Data Mutation Detected"},
            },
            {
                "type": "section",
                "text": {
                    "type": "mrkdwn",
                    "text": f"*{len(mutations)} records* changed without version increment.",
                },
            },
            {
                "type": "section",
                "text": {
                    "type": "mrkdwn",
                    "text": "```" + "\n".join(
                        f"{m['symbol']} @ {m['timestamp']}"
                        for m in mutations[:5]
                    ) + "```",
                },
            },
        ]

        payload = {"blocks": blocks}
        try:
            requests.post(self.webhook_url, json=payload, timeout=5)
        except Exception as e:
            logger.error(f"Failed to send Slack alert: {e}")


# ──────────────────────────────────────────────
# Scheduler: run at 4:15 PM ET daily
# ──────────────────────────────────────────────
if __name__ == "__main__":
    symbols = ["AAPL.US", "MSFT.US", "NVDA.US", "SPY.US"]  # Configure as needed

    job = DataSyncJob(symbols=symbols, interval="5m")

    # Run once on startup for immediate validation
    job.run()

    # Schedule daily runs
    schedule.every().day.at("16:15").do(job.run)

    while True:
        schedule.run_pending()
        time.sleep(60)

Key Production Considerations

The code above handles the common cases, but production deployments require additional attention:

  • Timezone handling: Always store timestamps in UTC (TIMESTAMPTZ in PostgreSQL). Convert to market time only at the display layer. Mixing timezones is the single most common cause of "missing data" reports that turn out to be timezone misalignments.
  • Connection pooling: If syncing hundreds of symbols, use a connection pool (psycopg2.pool.ThreadedConnectionPool) rather than opening a new connection per symbol.
  • Idempotency: The entire sync job should be idempotent — running it twice should produce the same result as running it once. UPSERT with version checks ensures this.
  • Monitoring: Track sync duration, record counts, and mutation rates over time. A sudden spike in mutations may indicate an upstream data source issue, not a backfill.

TickDB Integration: Leveraging Native Capabilities

TickDB's API provides several features that simplify these synchronization patterns. The /kline/latest endpoint gives you the current bar without needing to manage a "last sync timestamp" in your application state — useful for monitoring dashboards. The depth channel provides order book snapshots for real-time applications where batch synchronization is insufficient.

For the synchronization patterns described in this article, the key TickDB endpoints are:

  • GET /v1/market/kline — Historical OHLCV bars with version numbers for backfill detection.
  • GET /v1/market/kline/latest — Current bar for live monitoring.
  • GET /v1/symbols/available — Validate symbol availability before attempting a fetch.
# Validate symbols before running a sync
def validate_symbols(symbols: list[str]) -> dict[str, bool]:
    """Check which symbols are available via TickDB."""
    url = "https://api.tickdb.ai/v1/symbols/available"
    headers = {"X-API-Key": TICKDB_API_KEY}

    response = requests.get(url, headers=headers, timeout=(3.05, 10))
    response.raise_for_status()
    available = set(response.json().get("data", []))

    return {symbol: symbol in available for symbol in symbols}

Summary: A Layered Defense Against Data Corruption

Data synchronization is not a solved problem. Every pipeline will encounter duplicates, backfills, and silent mutations if it runs long enough. The solution is not a single clever trick but a layered defense:

  1. UPSERT with a composite primary key eliminates duplicate insertion at the database level. This is your first line of defense and should be non-negotiable.

  2. Version number tracking detects backfilled corrections from the data source. A 7-day lookback window catches most corporate action adjustments without excessive API usage.

  3. Hash validation catches silent mutations that version metadata might miss. Use it as an audit layer, not a primary mechanism.

  4. Mutation alerting closes the loop. A Slack alert for detected mutations transforms your sync job from a passive process into an active monitoring system.

  5. Idempotent design ensures that running the pipeline twice is safe. UPSERT with version checks is the foundation; do not introduce operations that violate this property.

The discipline is in the details. A 3% record loss rate will not show up in a cursory inspection of your database — it will show up in a backtest that outperforms live trading by a suspiciously clean margin. Build the sync pipeline right the first time, and you will never have to debug a silent data integrity failure at 2 AM.


Next Steps

If you are building a data pipeline for systematic trading, the patterns in this article give you a production-ready foundation. Sign up at tickdb.ai to get a free API key and start pulling clean, versioned OHLCV data.

If you are running a backtesting workflow and need 10+ years of historical data for cross-cycle validation, reach out to [email protected] for institutional data plans.

If you are interested in the depth channel for real-time order book monitoring, see the companion article on depth snapshot handling for event-driven strategies.


This article does not constitute investment advice. Market data synchronization and backtesting involve technical complexity; validate all pipelines with out-of-sample data before live deployment.