Drücken Sie ESC, um zu schließen

Krypto Liquidationskarte: So Traded man Zonen

Wo Smart Money die Retail-Tragödie ausnutzt: Praxishandbuch Liquidation Heatmap

Eine Liquidation Heatmap ist der Röntgenblick für die Hoffnung der Masse. Während Retail-Trader brav nach Lehrbuch-Charttechnik ihren Stop-Loss knapp hinter das letzte lokale Hoch oder Tief platzieren, sieht der Big Player auf der Map etwas völlig anderes: gebündelten Treibstoff für seinen nächsten Move.

Ein Market Maker kann in einem dünnen Markt schlichtweg keine Positionen im zweistelligen Millionenbereich aufbauen oder schließen, ohne sich selbst den Preis zu zerlegen. Er braucht Liquidität. Und diese Liquidität liefern die Margin Calls der anderen.

Die Short-Falle: Anatomie eines Short Squeezes

Stell dir folgende Lage vor: Ein Asset fällt seit drei Tagen durch. Die Masse riecht Schwäche, hebelt sich mit 20x bis 50x in Short-Positionen und knallt die Stop-Losses (oder lässt die Position bis zur Liquidation laufen) knapp hinter den nächsten Widerstand — etwa eine runde Marke oder die Oberkante einer Konsolidierung.

Was macht jetzt der Big Player, der eine fette Long-Position aufbauen oder seine Bestände im Top abladen will?

  • Preis komprimieren: Der Preis schleicht langsam nach oben. Keine Panik bei den Shorties, aber die Liquidationszone auf der Heatmap rückt in greifbare Nähe.
  • Der Impuls-Pikser: Ein gezielter Push katapultiert den Preis schlagartig 1,5 bis 2 % über den Widerstand.
  • Die Kettenreaktion: Die ersten Stops fliegen, hochgehebelte Positionen werden zwangsliquidiert. Die Börse haut automatisch Market-Buy-Orders raus, um die Shorts glattzustellen.
  • Volumen-Explosion: Die zwangsläufigen Market-Käufe der Börse knallen direkt in die Limit-Ask-Orders des Big Players. Er hat seine Sell-Orders genau dort platziert, wo der Panik-Demand der Short-Eindeckungen landet.

Der Retail-Trader denkt derweil: „Ausbruch! Trendwende, nix wie rein in Long!“. Er kauft mitten in den Hype-Impuls hinein. Genau in diesem Moment dreht der Big Player die Hand um — und der Kurs rauscht wie ein Stein nach unten, um direkt die nächsten Long-Stops einzusammeln.

Warum das Orderbuch kollabiert

Wenn die Liquidation eines 500.000-Dollar-Trades mit 20x Hebel feuert, haut die Börse eine unlimitierte Market Order raus. Stehen im Orderbuch auf der Gegenseite nicht genug Limit-Orders, schießt der Preis wie durch Butter durch die Level.

  • Limit-Abverkauf: Die Liquidation-Engine frisst sich durch das Orderbuch. Stehen bei 60.000 $ nur 5 BTC im Buch, aber 20 BTC müssen zwangsverkauft werden, bricht der Preis durch und rennt weiter, bis er ausreichend Tiefe findet.
  • Kaskadeneffekt: Jedes durchbrochene Level kaskadiert in die nächsten Stop-Losses und Liquidationen von Positionen mit kleinerem Hebel (z.B. den 10x-Jungs).
  • V-Shape Reversal (Dead Cat Bounce): Sobald die Kaskade die komplette Liquidität des Impulses verbrannt hat, ist das Orderbuch trocken. Der Market Maker, der exakt auf diesen Erschöpfungsmoment gewartet hat, greift zu Spottpreisen mit fetten Limit-Orders ab und drückt den Kurs instant zurück.

Trading Setups: Wie du Liquidationszonen profitabel tradest

Versuche niemals, blind mitten in das Liquidations-Inferno reinzutraden — der Spread und das Slippage fressen dich auf. Du tradest die Reaktion des Marktes, sobald das Feuer gelöscht ist.

Setup #1: Stop Hunting (Liquidity Sweep Trade)

  • Setup: Auf der Liquidation Map siehst du einen dichten Cluster von Shorts oder Longs, der weniger als 1 % vom aktuellen Kurs entfernt liegt.
  • Plan: Warte den Wick über das Level ab. Auf keinen Fall beim Ausbruch selbst einsteigen!
  • Execution: Der Preis schießt über das Niveau, das Volumen auf dem 1m-Chart explodiert, das Open Interest (OI) bricht ein (alle wurden rasiert). Sobald die Kerze mit einem langen Docht (Wick) hinter das Niveau zurückschließt — Feuer frei in die Gegenrichtung.
  • Stop-Loss: Messerscharf knapp über das Extremum des Impuls-Dochts setzen. Wird der Preis erneut durch diese Zone gedrückt, war es kein Liquidity Sweep, sondern ein echter Breakout. Dein Risiko muss knallhart begrenzt sein.

Monitoring automatisieren: Python-Script für Anomalien

Um nicht 24/7 vor den Monitoren zu kleben und auf den nächsten Kaskadeneffekt zu warten, hilft das folgende Script. Es überwachst den Liquidation-Feed der Börse in Echtzeit und schmeißt unübliche Zwangsschließungen direkt in die Konsole.

Pure Python, keine externen Heavy-Weight-Bibliotheken, läuft extrem stabil und verkraftet auch mal einen Reconnect ohne abzustürzen.

import json
import time
import threading
import logging
from collections import deque, OrderedDict
from statistics import median
from typing import Dict, Tuple

import requests
import websocket

from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry


class ProductionMarketEngine:
    """
    Production-grade Binance Futures liquidation / market-state engine.

    Data sources:
        - Binance Futures forceOrder WebSocket
        - Binance Futures ticker WebSocket
        - Binance Futures Open Interest REST

    Main analysis window:
        5 minutes

    Features:
        - Rolling liquidation window
        - Calendar-aligned 5-minute liquidation buckets
        - TTL event deduplication
        - Real market price feed
        - Exact 5-minute OI delta
        - Liquidation imbalance
        - Liquidation intensity
        - Price momentum
        - Volatility
        - Market-state classification
        - Signal scoring
        - WebSocket reconnect with exponential backoff
        - HTTP retries for Open Interest
        - Graceful shutdown
        - Thread-safe state
        - Runtime health monitoring
    """

    def __init__(
        self,
        symbol: str = "btcusdt",
        window_seconds: int = 300,
        oi_poll_seconds: int = 10,
        baseline_buckets: int = 288,
        min_baseline_buckets: int = 12,
    ):
        self.symbol = symbol.lower()
        self.symbol_upper = symbol.upper()

        self.window_seconds = window_seconds
        self.oi_poll_seconds = oi_poll_seconds

        # 288 × 5 minutes = 24 hours.
        self.baseline_buckets = baseline_buckets

        # Minimum completed buckets before baseline becomes reliable.
        self.min_baseline_buckets = min_baseline_buckets

        # --------------------------------------------------------------
        # Thread synchronization
        # --------------------------------------------------------------

        self.lock = threading.RLock()

        # Event used for interruptible thread shutdown / waiting.
        self.stop_event = threading.Event()

        # --------------------------------------------------------------
        # Runtime state
        # --------------------------------------------------------------

        self.is_running = False

        self.ws_combined = None
        self.ws_thread = None
        self.oi_thread = None
        self.maintenance_thread = None

        # --------------------------------------------------------------
        # Liquidation rolling window
        # --------------------------------------------------------------

        self.events = deque()

        # event_id -> received timestamp
        self.seen_ids = OrderedDict()

        self.dedup_ttl_seconds = max(
            window_seconds * 2,
            600,
        )

        # --------------------------------------------------------------
        # Real market price history
        # --------------------------------------------------------------

        self.price_history = deque()

        self.current_price = 0.0
        self.last_price_timestamp = 0.0

        # --------------------------------------------------------------
        # Open Interest
        # --------------------------------------------------------------

        self.oi_history = deque(
            maxlen=180
        )

        self.current_oi = 0.0
        self.last_oi_timestamp = 0.0

        # --------------------------------------------------------------
        # Calendar-aligned liquidation buckets
        # --------------------------------------------------------------

        # {
        #     bucket_timestamp: liquidation_notional
        # }
        self.bucket_accum = OrderedDict()

        # --------------------------------------------------------------
        # WebSocket state
        # --------------------------------------------------------------

        self.ws_connected = False
        self.last_ws_message = 0.0
        self.ws_reconnects = 0

        # --------------------------------------------------------------
        # Diagnostics
        # --------------------------------------------------------------

        self.total_liquidation_events = 0
        self.invalid_liquidation_events = 0

        # --------------------------------------------------------------
        # Logger
        # --------------------------------------------------------------

        self.logger = logging.getLogger(
            f"ProductionMarketEngine.{self.symbol_upper}"
        )

    # ==================================================================
    # TIME / BUCKET HELPERS
    # ==================================================================

    @staticmethod
    def _bucket_timestamp(timestamp: float) -> int:
        """
        Convert timestamp to calendar-aligned 5-minute bucket.
        """

        return int(timestamp // 300) * 300

    # ==================================================================
    # CLEANUP
    # ==================================================================

    def _clean_ttl_cache(self, now: float) -> None:
        """
        Remove expired liquidation event IDs.
        """

        cutoff = now - self.dedup_ttl_seconds

        while self.seen_ids:
            first_key, first_timestamp = next(
                iter(self.seen_ids.items())
            )

            if first_timestamp < cutoff:
                self.seen_ids.popitem(
                    last=False
                )
            else:
                break

    def _clean_old_events(self, now: float) -> None:
        """
        Keep liquidation and price data only inside
        the rolling analysis window.
        """

        cutoff = now - self.window_seconds

        while self.events:
            if self.events[0]["timestamp"] < cutoff:
                self.events.popleft()
            else:
                break

        while self.price_history:
            if self.price_history[0]["timestamp"] < cutoff:
                self.price_history.popleft()
            else:
                break

    def _clean_old_buckets(self, now: float) -> None:
        """
        Remove old liquidation buckets.

        This method is called from the dedicated maintenance thread,
        so bucket cleanup does NOT depend on new liquidation events.
        """

        current_bucket = self._bucket_timestamp(now)

        minimum_bucket = (
            current_bucket
            - (self.baseline_buckets + 2) * 300
        )

        while self.bucket_accum:
            first_bucket = next(
                iter(self.bucket_accum)
            )

            if first_bucket < minimum_bucket:
                self.bucket_accum.popitem(
                    last=False
                )
            else:
                break

    # ==================================================================
    # LIQUIDATION INGESTION
    # ==================================================================

    def _process_liquidation(
        self,
        order_data: dict,
    ) -> None:
        """
        Process Binance Futures forceOrder event.
        """

        if not isinstance(
            order_data,
            dict,
        ):
            return

        now = time.time()

        raw_side = str(
            order_data.get(
                "S",
                "",
            )
        ).upper()

        # --------------------------------------------------------------
        # Validate side before touching deduplication cache.
        # --------------------------------------------------------------

        if raw_side == "SELL":
            side = "LONG_LIQ"

        elif raw_side == "BUY":
            side = "SHORT_LIQ"

        else:
            with self.lock:
                self.invalid_liquidation_events += 1

            return

        try:
            execution_timestamp_ms = int(
                order_data.get(
                    "T",
                    int(now * 1000),
                )
            )

            event_timestamp = (
                execution_timestamp_ms
                / 1000.0
            )

            price = float(
                order_data.get(
                    "ap",
                    order_data.get(
                        "p",
                        0,
                    ),
                )
            )

            qty = float(
                order_data.get(
                    "q",
                    0,
                )
            )

            order_id = str(
                order_data.get(
                    "i",
                    "",
                )
            )

        except (
            TypeError,
            ValueError,
        ):
            with self.lock:
                self.invalid_liquidation_events += 1

            self.logger.warning(
                "Invalid liquidation event received."
            )

            return

        if price <= 0 or qty <= 0:
            with self.lock:
                self.invalid_liquidation_events += 1

            return

        notional = price * qty

        if notional <= 0:
            return

        # --------------------------------------------------------------
        # Event ID
        #
        # Prefer Binance order ID when available.
        # If unavailable, use deterministic composite fallback.
        # --------------------------------------------------------------

        if order_id:
            event_id = (
                f"{execution_timestamp_ms}_"
                f"{order_id}"
            )
        else:
            event_id = (
                f"{execution_timestamp_ms}_"
                f"{side}_"
                f"{price:.12f}_"
                f"{qty:.12f}"
            )

        event = {
            "id": event_id,
            "timestamp": event_timestamp,
            "side": side,
            "price": price,
            "qty": qty,
            "notional": notional,
        }

        with self.lock:
            self._clean_ttl_cache(
                now
            )

            if event_id in self.seen_ids:
                return

            self.seen_ids[event_id] = now

            self.events.append(
                event
            )

            # Calendar-aligned bucket.
            bucket = self._bucket_timestamp(
                event_timestamp
            )

            if bucket not in self.bucket_accum:
                self.bucket_accum[bucket] = 0.0

            self.bucket_accum[bucket] += (
                notional
            )

            self.total_liquidation_events += 1

    # ==================================================================
    # MARKET PRICE INGESTION
    # ==================================================================

    def _process_ticker(
        self,
        ticker_data: dict,
    ) -> None:
        """
        Process Binance ticker event.
        """

        try:
            price = float(
                ticker_data.get(
                    "c",
                    0,
                )
            )

        except (
            TypeError,
            ValueError,
        ):
            return

        if price <= 0:
            return

        now = time.time()

        with self.lock:
            self.current_price = price
            self.last_price_timestamp = now

            self.price_history.append(
                {
                    "timestamp": now,
                    "price": price,
                }
            )

    # ==================================================================
    # OPEN INTEREST HTTP SESSION
    # ==================================================================

    @staticmethod
    def _create_http_session() -> requests.Session:
        """
        Create requests Session with connection pooling
        and retry policy.
        """

        session = requests.Session()

        retry = Retry(
            total=4,
            connect=4,
            read=4,
            status=4,

            backoff_factor=0.5,

            status_forcelist=(
                429,
                500,
                502,
                503,
                504,
            ),

            allowed_methods=frozenset(
                [
                    "GET",
                ]
            ),

            respect_retry_after_header=True,
        )

        adapter = HTTPAdapter(
            max_retries=retry,
            pool_connections=4,
            pool_maxsize=4,
        )

        session.mount(
            "https://",
            adapter,
        )

        session.headers.update(
            {
                "User-Agent":
                    "ProductionMarketEngine/1.0"
            }
        )

        return session

    # ==================================================================
    # OPEN INTEREST POLLER
    # ==================================================================

    def _poll_open_interest(
        self,
    ) -> None:
        """
        Poll Binance Futures Open Interest.
        """

        url = (
            "https://fapi.binance.com"
            "/fapi/v1/openInterest"
        )

        session = (
            self._create_http_session()
        )

        try:
            while not self.stop_event.is_set():

                try:
                    response = session.get(
                        url,
                        params={
                            "symbol":
                                self.symbol_upper
                        },
                        timeout=5,
                    )

                    response.raise_for_status()

                    data = response.json()

                    oi = float(
                        data.get(
                            "openInterest",
                            0,
                        )
                    )

                    if oi <= 0:
                        raise ValueError(
                            "Invalid Open Interest."
                        )

                    now = time.time()

                    with self.lock:
                        self.current_oi = oi
                        self.last_oi_timestamp = now

                        self.oi_history.append(
                            {
                                "timestamp": now,
                                "oi": oi,
                            }
                        )

                except requests.RequestException as exc:
                    self.logger.warning(
                        "Open Interest request failed: %s",
                        exc,
                    )

                except (
                    ValueError,
                    TypeError,
                ) as exc:
                    self.logger.warning(
                        "Invalid Open Interest response: %s",
                        exc,
                    )

                except Exception:
                    self.logger.exception(
                        "Unexpected Open Interest error."
                    )

                # Interruptible wait.
                self.stop_event.wait(
                    self.oi_poll_seconds
                )

        finally:
            session.close()

    # ==================================================================
    # MAINTENANCE THREAD
    # ==================================================================

    def _maintenance_loop(
        self,
    ) -> None:
        """
        Periodic memory/state cleanup.

        Important:
        cleanup is independent of liquidation activity.
        """

        while not self.stop_event.wait(10):
            now = time.time()

            with self.lock:
                self._clean_ttl_cache(
                    now
                )

                self._clean_old_events(
                    now
                )

                self._clean_old_buckets(
                    now
                )

    # ==================================================================
    # OI 5-MINUTE DELTA
    # ==================================================================

    def _get_oi_change_5m(
        self,
        now: float,
    ) -> Tuple[float, bool]:
        """
        Calculate actual 5-minute Open Interest change.

        Returns:
            change_pct
            data_ready
        """

        if len(self.oi_history) < 2:
            return 0.0, False

        cutoff = (
            now
            - self.window_seconds
        )

        baseline_sample = None

        for sample in self.oi_history:

            if sample["timestamp"] <= cutoff:
                baseline_sample = sample

            else:
                break

        # Not enough historical data yet.
        if baseline_sample is None:
            oldest = self.oi_history[0]

            if (
                now
                - oldest["timestamp"]
                < self.window_seconds * 0.8
            ):
                return 0.0, False

            baseline_sample = oldest

        latest_sample = (
            self.oi_history[-1]
        )

        first_oi = baseline_sample["oi"]
        last_oi = latest_sample["oi"]

        if first_oi <= 0:
            return 0.0, False

        change_pct = (
            (
                last_oi
                - first_oi
            )
            / first_oi
        ) * 100.0

        return change_pct, True

    # ==================================================================
    # PRICE METRICS
    # ==================================================================

    def _get_price_metrics(
        self,
    ) -> Tuple[
        float,
        float,
        float,
        bool,
    ]:
        """
        Returns:

            current price
            5m price change
            5m high-low volatility
            ready
        """

        if len(
            self.price_history
        ) < 2:
            return (
                0.0,
                0.0,
                0.0,
                False,
            )

        first_price = (
            self.price_history[0]["price"]
        )

        last_price = (
            self.price_history[-1]["price"]
        )

        if first_price <= 0:
            return (
                0.0,
                0.0,
                0.0,
                False,
            )

        price_change_pct = (
            (
                last_price
                - first_price
            )
            / first_price
        ) * 100.0

        prices = [
            item["price"]
            for item in self.price_history
        ]

        high_price = max(prices)
        low_price = min(prices)

        if low_price > 0:
            volatility_pct = (
                (
                    high_price
                    - low_price
                )
                / low_price
            ) * 100.0

        else:
            volatility_pct = 0.0

        return (
            last_price,
            price_change_pct,
            volatility_pct,
            True,
        )

    # ==================================================================
    # LIQUIDATION METRICS
    # ==================================================================

    def _get_liquidation_metrics(
        self,
    ) -> Tuple[
        float,
        float,
        float,
        float,
    ]:
        """
        Returns:

            long liquidation
            short liquidation
            total liquidation
            imbalance
        """

        long_liq = 0.0
        short_liq = 0.0

        for event in self.events:

            if event["side"] == "LONG_LIQ":
                long_liq += event["notional"]

            elif event["side"] == "SHORT_LIQ":
                short_liq += event["notional"]

        total_liq = (
            long_liq
            + short_liq
        )

        if total_liq > 0:
            imbalance = (
                (
                    long_liq
                    - short_liq
                )
                / total_liq
            )

        else:
            imbalance = 0.0

        return (
            long_liq,
            short_liq,
            total_liq,
            imbalance,
        )

    # ==================================================================
    # BASELINE
    # ==================================================================

    def _get_baseline(
        self,
    ) -> Tuple[
        float,
        int,
        bool,
    ]:
        """
        Calculate median liquidation volume
        from completed calendar-aligned 5m buckets.

        Current incomplete bucket is excluded.
        """

        now = time.time()

        current_bucket = (
            self._bucket_timestamp(
                now
            )
        )

        completed = []

        for (
            bucket_timestamp,
            volume,
        ) in self.bucket_accum.items():

            if (
                bucket_timestamp
                < current_bucket
            ):
                completed.append(
                    volume
                )

        if not completed:
            return (
                0.0,
                0,
                False,
            )

        completed = completed[
            -self.baseline_buckets:
        ]

        baseline = median(
            completed
        )

        reliable = (
            len(completed)
            >= self.min_baseline_buckets
        )

        return (
            float(baseline),
            len(completed),
            reliable,
        )

    # ==================================================================
    # UTILITIES
    # ==================================================================

    @staticmethod
    def _clamp(
        value: float,
        minimum: float,
        maximum: float,
    ) -> float:
        return max(
            minimum,
            min(
                maximum,
                value,
            ),
        )

    # ==================================================================
    # SIGNAL SCORES
    # ==================================================================

    def _calculate_scores(
        self,
        intensity: float,
        imbalance: float,
        price_change_pct: float,
        oi_change_pct: float,
        oi_ready: bool,
    ) -> Dict[str, float]:
        """
        Calculate normalized signal-strength scores.

        Scores are signal strengths, not probabilities.
        """

        liquidation_pressure = (
            self._clamp(
                (
                    intensity
                    / 5.0
                ) * 100.0,
                0.0,
                100.0,
            )
        )

        imbalance_score = (
            abs(imbalance)
            * 100.0
        )

        price_impulse = (
            self._clamp(
                (
                    abs(
                        price_change_pct
                    )
                    / 3.0
                )
                * 100.0,
                0.0,
                100.0,
            )
        )

        if oi_ready:
            oi_contraction = (
                self._clamp(
                    (
                        abs(
                            min(
                                oi_change_pct,
                                0.0,
                            )
                        )
                        / 3.0
                    )
                    * 100.0,
                    0.0,
                    100.0,
                )
            )

        else:
            oi_contraction = 0.0

        return {
            "liquidation_pressure":
                liquidation_pressure,

            "imbalance":
                imbalance_score,

            "price_impulse":
                price_impulse,

            "oi_contraction":
                oi_contraction,
        }

    # ==================================================================
    # MARKET ANALYSIS
    # ==================================================================

    def analyze_market(
        self,
    ) -> Tuple[str, Dict]:
        """
        Main market-state classifier.
        """

        now = time.time()

        with self.lock:

            self._clean_old_events(
                now
            )

            self._clean_ttl_cache(
                now
            )

            self._clean_old_buckets(
                now
            )

            (
                last_price,
                price_change_pct,
                volatility_pct,
                price_ready,
            ) = self._get_price_metrics()

            if not price_ready:
                return (
                    "INITIALIZING_PRICE_FEED",
                    {},
                )

            (
                oi_change_pct,
                oi_ready,
            ) = self._get_oi_change_5m(
                now
            )

            (
                long_liq,
                short_liq,
                total_liq,
                imbalance,
            ) = self._get_liquidation_metrics()

            (
                baseline_med,
                baseline_count,
                baseline_ready,
            ) = self._get_baseline()

            # ----------------------------------------------------------
            # No liquidation activity.
            # ----------------------------------------------------------

            if total_liq <= 0:

                metrics = {
                    "price":
                        last_price,

                    "price_change_5m":
                        price_change_pct,

                    "volatility_5m":
                        volatility_pct,

                    "oi_change_5m":
                        oi_change_pct,

                    "oi_ready":
                        oi_ready,

                    "total_liq_usd":
                        0.0,

                    "long_liq_usd":
                        0.0,

                    "short_liq_usd":
                        0.0,

                    "imbalance":
                        0.0,

                    "intensity":
                        0.0,

                    "baseline_med":
                        baseline_med,

                    "baseline_buckets":
                        baseline_count,

                    "baseline_ready":
                        baseline_ready,

                    "scores": {
                        "liquidation_pressure":
                            0.0,

                        "imbalance":
                            0.0,

                        "price_impulse":
                            self._clamp(
                                (
                                    abs(
                                        price_change_pct
                                    )
                                    / 3.0
                                )
                                * 100.0,
                                0.0,
                                100.0,
                            ),

                        "oi_contraction":
                            0.0,
                    },

                    "confidence":
                        0.0,
                }

                if (
                    abs(
                        price_change_pct
                    ) < 0.3
                ):
                    return (
                        "NO_LIQUIDATION_ACTIVITY",
                        metrics,
                    )

                return (
                    "PRICE_MOVEMENT_NO_LIQUIDATION",
                    metrics,
                )

            # ----------------------------------------------------------
            # Intensity.
            # ----------------------------------------------------------

            if (
                baseline_ready
                and baseline_med > 0
            ):
                intensity = (
                    total_liq
                    / baseline_med
                )

            else:
                intensity = 0.0

            # ----------------------------------------------------------
            # Scores.
            # ----------------------------------------------------------

            scores = (
                self._calculate_scores(
                    intensity=intensity,
                    imbalance=imbalance,
                    price_change_pct=price_change_pct,
                    oi_change_pct=oi_change_pct,
                    oi_ready=oi_ready,
                )
            )

            # ----------------------------------------------------------
            # Confidence.
            # ----------------------------------------------------------

            confidence_components = [
                scores[
                    "liquidation_pressure"
                ],

                scores[
                    "imbalance"
                ],

                scores[
                    "price_impulse"
                ],
            ]

            if oi_ready:
                confidence_components.append(
                    scores[
                        "oi_contraction"
                    ]
                )

            confidence = (
                sum(
                    confidence_components
                )
                / len(
                    confidence_components
                )
            )

            metrics = {
                "price":
                    last_price,

                "price_change_5m":
                    price_change_pct,

                "volatility_5m":
                    volatility_pct,

                "oi_change_5m":
                    oi_change_pct,

                "oi_ready":
                    oi_ready,

                "total_liq_usd":
                    total_liq,

                "long_liq_usd":
                    long_liq,

                "short_liq_usd":
                    short_liq,

                "imbalance":
                    imbalance,

                "intensity":
                    intensity,

                "baseline_med":
                    baseline_med,

                "baseline_buckets":
                    baseline_count,

                "baseline_ready":
                    baseline_ready,

                "scores":
                    scores,

                "confidence":
                    confidence,
            }

            # ==========================================================
            # CASCADES
            # ==========================================================

            if (
                baseline_ready
                and intensity >= 3.0
                and imbalance >= 0.60
                and price_change_pct <= -2.0
            ):
                return (
                    "CASCADING_DOWNSELL",
                    metrics,
                )

            if (
                baseline_ready
                and intensity >= 3.0
                and imbalance <= -0.60
                and price_change_pct >= 2.0
            ):
                return (
                    "CASCADING_PUMP",
                    metrics,
                )

            # ==========================================================
            # LONG CAPITULATION
            # ==========================================================

            if (
                baseline_ready
                and oi_ready
                and intensity >= 2.0
                and imbalance >= 0.50
                and price_change_pct <= -0.50
                and oi_change_pct <= -0.50
            ):
                return (
                    "LONG_CAPITULATION",
                    metrics,
                )

            # ==========================================================
            # SHORT SQUEEZE EXHAUSTION
            # ==========================================================

            if (
                baseline_ready
                and oi_ready
                and intensity >= 2.0
                and imbalance <= -0.50
                and price_change_pct >= 0.50
                and oi_change_pct <= -0.50
            ):
                return (
                    "SHORT_SQUEEZE_EXHAUSTION",
                    metrics,
                )

            # ==========================================================
            # DELEVERAGING
            # ==========================================================

            if (
                oi_ready
                and imbalance >= 0.50
                and oi_change_pct <= -0.30
            ):
                return (
                    "LONG_DELEVERAGING",
                    metrics,
                )

            if (
                oi_ready
                and imbalance <= -0.50
                and oi_change_pct <= -0.30
            ):
                return (
                    "SHORT_DELEVERAGING",
                    metrics,
                )

            # ==========================================================
            # EXTREME LIQUIDATION SPIKES
            # ==========================================================

            if (
                baseline_ready
                and intensity >= 2.5
                and imbalance >= 0.60
            ):
                return (
                    "EXTREME_LONG_LIQUIDATION_SPIKE",
                    metrics,
                )

            if (
                baseline_ready
                and intensity >= 2.5
                and imbalance <= -0.60
            ):
                return (
                    "EXTREME_SHORT_LIQUIDATION_SPIKE",
                    metrics,
                )

            # ==========================================================
            # ELEVATED LIQUIDATION
            # ==========================================================

            if imbalance >= 0.40:
                return (
                    "ELEVATED_LONG_LIQUIDATION",
                    metrics,
                )

            if imbalance <= -0.40:
                return (
                    "ELEVATED_SHORT_LIQUIDATION",
                    metrics,
                )

            # ==========================================================
            # BASELINE WARMUP
            # ==========================================================

            if not baseline_ready:
                return (
                    "WARMING_BASELINE",
                    metrics,
                )

            if intensity >= 1.5:
                return (
                    "BALANCED_HIGH_LIQUIDATION_VOLUME",
                    metrics,
                )

            return (
                "BALANCED",
                metrics,
            )

    # ==================================================================
    # WEBSOCKET CREATION
    # ==================================================================

    def _create_websocket(self):
        """
        Create Binance combined WebSocket.
        """

        combined_stream = (
            "wss://fstream.binance.com/stream"
            f"?streams="
            f"{self.symbol}@forceOrder/"
            f"{self.symbol}@ticker"
        )

        def on_open(ws):
            with self.lock:
                self.ws_connected = True

            self.logger.info(
                "WebSocket connected: %s",
                self.symbol_upper,
            )

        def on_message(
            ws,
            message,
        ):
            self.last_ws_message = (
                time.time()
            )

            try:
                payload = json.loads(
                    message
                )

                stream = payload.get(
                    "stream",
                    "",
                )

                data = payload.get(
                    "data",
                    {},
                )

                if (
                    "forceOrder"
                    in stream
                ):
                    self._process_liquidation(
                        data.get(
                            "o",
                            {},
                        )
                    )

                elif (
                    "ticker"
                    in stream
                ):
                    self._process_ticker(
                        data
                    )

            except json.JSONDecodeError:
                self.logger.warning(
                    "Invalid WebSocket JSON."
                )

            except Exception:
                self.logger.exception(
                    "Unexpected WebSocket message error."
                )

        def on_error(
            ws,
            error,
        ):
            with self.lock:
                self.ws_connected = False

            self.logger.warning(
                "WebSocket error: %s",
                error,
            )

        def on_close(
            ws,
            close_status_code,
            close_msg,
        ):
            with self.lock:
                self.ws_connected = False

            self.logger.warning(
                "WebSocket closed: "
                "code=%s message=%s",
                close_status_code,
                close_msg,
            )

        return websocket.WebSocketApp(
            combined_stream,
            on_open=on_open,
            on_message=on_message,
            on_error=on_error,
            on_close=on_close,
        )

    # ==================================================================
    # WEBSOCKET LOOP
    # ==================================================================

    def _websocket_loop(
        self,
    ) -> None:
        """
        Persistent WebSocket loop.

        Uses exponential reconnect backoff capped at 60 seconds.
        Shutdown is interruptible through stop_event.
        """

        backoff = 1

        while not self.stop_event.is_set():

            try:
                self.ws_combined = (
                    self._create_websocket()
                )

                self.ws_reconnects += 1

                self.logger.info(
                    "Starting WebSocket "
                    "connection attempt #%d.",
                    self.ws_reconnects,
                )

                self.ws_combined.run_forever(
                    ping_interval=30,
                    ping_timeout=10,
                    skip_utf8_validation=False,
                )

                # If run_forever returns normally after
                # an established connection, reset backoff.
                if self.is_running:
                    backoff = 1

            except Exception as exc:
                self.logger.error(
                    "WebSocket run_forever crashed: %s",
                    exc,
                )

            finally:
                with self.lock:
                    self.ws_connected = False

            if self.stop_event.is_set():
                break

            self.logger.info(
                "Reconnecting WebSocket in %d seconds...",
                backoff,
            )

            # Interruptible exponential backoff.
            if self.stop_event.wait(
                backoff
            ):
                break

            backoff = min(
                backoff * 2,
                60,
            )

        self.logger.info(
            "WebSocket loop stopped."
        )

    # ==================================================================
    # START
    # ==================================================================

    def start(
        self,
    ) -> None:
        """
        Start all background threads.
        """

        with self.lock:

            if self.is_running:
                self.logger.warning(
                    "Engine is already running."
                )
                return

            self.is_running = True

            self.stop_event.clear()

        # --------------------------------------------------------------
        # Open Interest thread
        # --------------------------------------------------------------

        self.oi_thread = threading.Thread(
            target=self._poll_open_interest,
            name=(
                f"OI_Poll_"
                f"{self.symbol_upper}"
            ),
            daemon=True,
        )

        self.oi_thread.start()

        # --------------------------------------------------------------
        # WebSocket thread
        # --------------------------------------------------------------

        self.ws_thread = threading.Thread(
            target=self._websocket_loop,
            name=(
                f"WS_Loop_"
                f"{self.symbol_upper}"
            ),
            daemon=True,
        )

        self.ws_thread.start()

        # --------------------------------------------------------------
        # Maintenance thread
        # --------------------------------------------------------------

        self.maintenance_thread = (
            threading.Thread(
                target=self._maintenance_loop,
                name=(
                    f"Maintenance_"
                    f"{self.symbol_upper}"
                ),
                daemon=True,
            )
        )

        self.maintenance_thread.start()

        self.logger.info(
            "ProductionMarketEngine started "
            "successfully: %s",
            self.symbol_upper,
        )

    # ==================================================================
    # STOP
    # ==================================================================

    def stop(
        self,
    ) -> None:
        """
        Graceful and interruptible shutdown.

        Order:
            1. Signal all workers to stop.
            2. Close WebSocket.
            3. Wait for threads.
            4. Reset runtime state.
        """

        with self.lock:

            if not self.is_running:
                return

            self.logger.info(
                "Stopping ProductionMarketEngine..."
            )

            self.is_running = False

            self.stop_event.set()

            self.ws_connected = False

        # --------------------------------------------------------------
        # Explicitly close WebSocket.
        # This unblocks run_forever().
        # --------------------------------------------------------------

        ws = self.ws_combined

        if ws is not None:

            try:
                ws.close()

            except Exception as exc:
                self.logger.debug(
                    "Error closing WebSocket: %s",
                    exc,
                )

        # --------------------------------------------------------------
        # Join worker threads.
        # --------------------------------------------------------------

        current_thread = (
            threading.current_thread()
        )

        threads = [
            self.ws_thread,
            self.oi_thread,
            self.maintenance_thread,
        ]

        for thread in threads:

            if (
                thread is not None
                and thread.is_alive()
                and thread is not current_thread
            ):
                thread.join(
                    timeout=5.0
                )

        with self.lock:
            self.ws_combined = None
            self.ws_thread = None
            self.oi_thread = None
            self.maintenance_thread = None

        self.logger.info(
            "ProductionMarketEngine stopped."
        )

    # ==================================================================
    # HEALTH
    # ==================================================================

    def health(
        self,
    ) -> Dict:
        """
        Runtime health information.
        """

        now = time.time()

        with self.lock:

            ws_age = (
                now
                - self.last_ws_message
                if self.last_ws_message > 0
                else None
            )

            oi_age = (
                now
                - self.last_oi_timestamp
                if self.last_oi_timestamp > 0
                else None
            )

            price_age = (
                now
                - self.last_price_timestamp
                if self.last_price_timestamp > 0
                else None
            )

            return {
                "running":
                    self.is_running,

                "symbol":
                    self.symbol_upper,

                "websocket_connected":
                    self.ws_connected,

                "websocket_last_message_age":
                    ws_age,

                "websocket_reconnects":
                    self.ws_reconnects,

                "price":
                    self.current_price,

                "price_age":
                    price_age,

                "open_interest":
                    self.current_oi,

                "oi_age":
                    oi_age,

                "liquidation_events":
                    self.total_liquidation_events,

                "invalid_liquidation_events":
                    self.invalid_liquidation_events,

                "dedup_cache_size":
                    len(
                        self.seen_ids
                    ),

                "rolling_liquidation_events":
                    len(
                        self.events
                    ),

                "price_history_size":
                    len(
                        self.price_history
                    ),

                "oi_history_size":
                    len(
                        self.oi_history
                    ),

                "baseline_buckets":
                    len(
                        self.bucket_accum
                    ),
            }


# ======================================================================
# CONSOLE OUTPUT
# ======================================================================

def print_market_state(
    state: str,
    metrics: Dict,
) -> None:
    """
    Human-readable realtime market state.
    """

    if not metrics:

        print(
            f"\r[{time.strftime('%H:%M:%S')}] "
            f"Инициализация потоков...",
            end="",
            flush=True,
        )

        return

    scores = metrics.get(
        "scores",
        {},
    )

    print(
        f"\n[{time.strftime('%H:%M:%S')}] "
        f"СТАТУС: >>> {state} <<<"
    )

    print(
        f" ├─ Цена BTC:             "
        f"${metrics['price']:,.2f}"
    )

    print(
        f" ├─ Цена 5M:              "
        f"{metrics['price_change_5m']:+.2f}%"
    )

    print(
        f" ├─ Волатильность 5M:     "
        f"{metrics['volatility_5m']:.2f}%"
    )

    oi_status = (
        "READY"
        if metrics["oi_ready"]
        else "WARMING"
    )

    print(
        f" ├─ Open Interest 5M:     "
        f"{metrics['oi_change_5m']:+.2f}% "
        f"[{oi_status}]"
    )

    print(
        f" ├─ Ликвидации 5M:        "
        f"${metrics['total_liq_usd']:,.0f}"
    )

    print(
        f" │   ├─ Long:             "
        f"${metrics['long_liq_usd']:,.0f}"
    )

    print(
        f" │   └─ Short:            "
        f"${metrics['short_liq_usd']:,.0f}"
    )

    print(
        f" ├─ Imbalance:            "
        f"{metrics['imbalance']:+.3f}"
    )

    print(
        f" ├─ Intensity:            "
        f"{metrics['intensity']:.2f}x"
    )

    print(
        f" ├─ Baseline median:      "
        f"${metrics['baseline_med']:,.0f}"
    )

    print(
        f" ├─ Baseline buckets:     "
        f"{metrics['baseline_buckets']}"
        f"{' [READY]' if metrics['baseline_ready'] else ' [WARMING]'}"
    )

    print(
        f" ├─ Liquidation score:    "
        f"{scores.get('liquidation_pressure', 0):.0f}/100"
    )

    print(
        f" ├─ Imbalance score:      "
        f"{scores.get('imbalance', 0):.0f}/100"
    )

    print(
        f" ├─ Price impulse:        "
        f"{scores.get('price_impulse', 0):.0f}/100"
    )

    print(
        f" ├─ OI contraction:       "
        f"{scores.get('oi_contraction', 0):.0f}/100"
    )

    print(
        f" └─ Confidence:           "
        f"{metrics.get('confidence', 0):.0f}/100"
    )


# ======================================================================
# LOGGING
# ======================================================================

def configure_logging() -> None:
    logging.basicConfig(
        level=logging.INFO,
        format=(
            "%(asctime)s | "
            "%(levelname)s | "
            "%(name)s | "
            "%(message)s"
        ),
        datefmt="%Y-%m-%d %H:%M:%S",
    )


# ======================================================================
# APPLICATION
# ======================================================================

if __name__ == "__main__":

    configure_logging()

    engine = ProductionMarketEngine(
        symbol="btcusdt",
        window_seconds=300,
        oi_poll_seconds=10,
        baseline_buckets=288,
        min_baseline_buckets=12,
    )

    print(
        "Starting Production Market Engine..."
    )

    print(
        "Streams:"
    )

    print(
        "  - Binance Futures forceOrder"
    )

    print(
        "  - Binance Futures ticker"
    )

    print(
        "  - Binance Futures Open Interest"
    )

    print(
        "Analysis window: 5 minutes"
    )

    print(
        "Baseline: 288 calendar-aligned 5M buckets"
    )

    print(
        "Baseline warm-up: "
        f"{engine.min_baseline_buckets} buckets"
    )

    engine.start()

    try:

        while True:

            state, metrics = (
                engine.analyze_market()
            )

            print_market_state(
                state,
                metrics,
            )

            time.sleep(10)

    except KeyboardInterrupt:

        print(
            "\n\nStopping engine..."
        )

    finally:

        engine.stop()

        print(
            "Engine stopped."
        )

So baust du das Script in deinen Workflow ein:

  • Passe den Schwellenwert (`threshold_usd`) an deine Depotgröße und das Asset an. Bei den Heavyweights ($BTC, $ETH) liefert ein Schwellenwert ab 500.000 $ verlässliche Signale für echte Whales.
  • Schlägt das Script an: Sofort den Footprint-Chart oder die Cluster prüfen. Liegt die Großliquidation an einem Schlüssel-Support oder -Resistance, ist das dein primäres Signal für die Vorbereitung eines Counter-Trend-Trades.

Warum die Heatmap Solotradern etwas vorgaukelt – und wie du sie in Kombination mit Funding Rates richtig liest

Der häufigste Anfängerfehler überhaupt: die Liquidations-Heatmap isoliert zu betrachten. Eine dichte Ansammlung roter oder grüner Balken im Chart ist absolut keine Garantie für einen Reversal. Die Map zeigt dir lediglich potenziellen Treibstoff, verrät dir aber nicht, wohin die eigentliche Maschine steuern will.

Um einen Fakeout von einem echten, gewaltigen Liquidity Sweep zu unterscheiden, musst du die Funding Rate und das Open Interest (OI) zusammensetzen:

  • Extrem positive Funding Rate + steigendes Open Interest: Die Masse ballert vollgepumpt mit Hebeln in Long-Positionen. Jeder rechnet fest mit dem Ausbruch nach oben. Das ist der perfekte Nährboden für einen gnadenlosen Long Squeeze nach unten. Die Liquidation Map auf der Unterseite wird völlig überladen sein.
  • Tief negative Funding Rate + steigendes Open Interest: Der Markt ist vollgestopft mit Shorties. Jeder kleine Impuls nach oben löst eine Kettenreaktion in Form eines Short Squeezes aus, weil die geliehenen Shorts panisch zurückgekauft werden müssen.

Siehst du auf der Map ein massives Cluster an Short-Liquidationen, während die Funding Rate über mehrere Stunden tief im Minus klebt, geht die Wahrscheinlichkeit für einen brutalen Push nach oben gen 100 %. Kein Big Player lässt sich diese feine Liquidität entgehen.

Setup #2: Der echte Ausbruch nach dem Liquidity Sweep (Trendfortsetzung)

Nicht jeder Abfischen von Liquidität endet in einem prompten V-Reversal oder einem langen Docht. Manchmal ist der aufgestaute Treibstoff so gewaltig, dass der Preis nicht nur einen kurzen Liquidation Spike hinlegt, sondern direkt in einen mehrere Stunden andauernden Impulstrend übergeht.

So erkennst du einen echten Breakout durch eine Liquidationszone:

  • Kein Pullback: Der Preis rasiert das Level mit den gebündelten Stops, hinterlässt aber anstelle eines langen Dochts eine stramme 1-Stunden-Kerze, die sauber ober- bzw. unterhalb schließt.
  • Verhalten des CVD (Cumulative Volume Delta): Das Delta der aggressiven Market-Buys/-Sells schießt zusammen mit dem Preis vertikal in die Höhe. Das bedeutet: Hier kaufen nicht nur die erzwungenen Liquidationen der Börse, sondern auch echte Market-Käufer schmeißen sich mit voller Kapazität in die Bewegung.
  • Open Interest nach dem Spike: Nachdem der Liquidationsschwall durch ist, bricht das OI nicht komplett ein (oder fängt sich nach einem kurzen Dipp sofort wieder). Das ist das eindeutige Zeichen dafür, dass frisch frei gewordene Liquidität von neuen Big Playern übernommen wird und nicht bloß alte Positionen glattgestellt werden.

In diesem Fall erfolgt der Einstieg beim Retest der durchbrochenen Zone mit einem knallharten Stop hinter der Level-Kante.

Eiserne Risk-Management-Regeln bei der Jagd nach Liquidität

Das Trading in Liquidationszonen ist ein Spaziergang durch ein Minenfeld. Wenn du das Timing um zwei Sekunden verhaust oder ohne Stop reingehst, wirst du ganz schnell selbst Teil der nächsten Heatmap.

  • Vergiss Hebel über 5x–10x: Wenn du gegen hochvolatile Impulse beim Stop-Fishing tradest, kann dir der Markt innerhalb von Millisekunden 1–2 % Slippage reindrücken. Ein hoher Hebel zerschießt dir dein Depot noch bevor dein eigener Stop-Loss überhaupt am Orderbuch ankommt.
  • Der Stop-Loss ist Gesetz: Hat der Preis das Liquidationslevel durchschlagen und zieht ohne dich durch – verbillige nicht! Akzeptiere deinen fixen Verlust (1–2 % vom Depowert), gestehe dem Markt recht zu und warte auf die nächste saubere Setup-Zone.
  • Ignoriere Mini-Cluster: Reagiere nicht auf jeden kleinen Pip im Orderbuch. Ein richtiger Big Player jagt nur das große Geld. Konzentriere dich auf Cluster, deren Volumen das Vielfache des normalen Minutenvolumens des Assets übersteigen.

Trading ist kein Erraten der nächsten grünen Kerze. Es ist reine Stochastik, Wahrscheinlichkeitsrechnung und das eiskalte Ausnutzen der Fehler anderer zu deinen Gunsten. Analysiere die Liquidity Map, denke wie ein Market Maker und lass deine Position niemals ohne Schutz laufen.

Diesen Blogbeitrag zusammenfassen mit:

FAQ

Eine Crypto Liquidation Map ist ein visuelles Daten-Tool, das die geschätzten Preisniveaus berechnet, bei denen gehebelte Perpetual-Futures-Positionen der Zwangsliquidierung unterliegen. Der Algorithmus ermittelt die Dichte erzwungener Marktorders, indem er das Open Interest mit verschiedenen Hebelstufen (25x, 50x, 100x) kombiniert, um hochliquide Zonen zu identifizieren.

Händler nutzen Liquidity-Sweep-Strategien, indem sie große Liquidations-Cluster als Ziele für Preisbewegungen identifizieren. Take-Profit-Orders werden direkt in die Kaskadenzone gesetzt, während Reversal-Einstiege erfolgen, sobald die erzwungenen Zwangverkäufe oder -käufe vollständig ausgeführt und von Limit-Orders im Orderbuch absorbiert wurden.

Liquidation Maps liefern statistische Wahrscheinlichkeitsmodelle statt exakter Daten aus dem Orderbuch der Börsen. Das Tool rekonstruiert die Hebelverteilung durch die kontinuierliche Aggregation von Echtzeit-Open-Interest, Funding Rates und Delta-Volumenströmen, was eine hohe mathematische Präzision bei der Bestimmung mikrostruktureller Unterstützungs- und Widerstandszonen ermöglicht.
Martyn Borkowski

I am a crypto trader specializing in digital assets and blockchain markets.

My focus is on identifying opportunities, managing risk, and optimizing strategies to achieve consistent growth in the fast-evolving world of cryptocurrency.

Verification & Professional Profiles: X Profile

...

Diskussion beitreten

Ihre E-Mail-Adresse wird nicht veröffentlicht. Erforderliche Felder sind markiert *