اضغط على ESC للإغلاق

خريطة تصفية العملات الرقمية: دليل التداول المتقدم

أين يُصفي «الأموال الذكية» المتداولين الصغار: الدليل العملي للتعامل مع خريطة التصفية

خريطة التصفية الحرارية (Liquidation Heatmap) هي بمثابة أشعة سينية لآمال الآخرين. بينما يضع متداول التجزئة (Retail Trader) أمر إيقاف الخسارة (Stop-Loss) بناءً على التحليل الفني خلف أقرب قمة أو قاع محلي، يرى الحوت الكبير تكتل هذه النقاط كوقود جاهز لدفع سعره.

صانع السوق (Market Maker) لا يجد عمقًا كافيًا لبناء أو إغلاق صفقات بمليارات الدولارات في سوق ضعيف السيولة. هو بحاجة إلى سيولة عارمة، وهذه السيولة تأتي تحديدًا من نداءات الهامش (Margin Calls) الخاصة بالآخرين.

مصيدة البائعين: ميكانيكية الـ Short Squeeze

تخيل هذا السيناريو: الأصل في حالة هبوط منذ ثلاثة أيام. المتداولون يرون ضعفًا في السوق، فيدخلون صفقات بيع (Short) برافعة مالية ضخمة من 20x إلى 50x ويضعون أوامر وقف الخسارة (أو يحتفظون بالهامش حتى التصفية) خلف أقرب مستوى مقاومة محلي — على سبيل المثال، عند مستوى رقمي صحيح أو عند الحد الأعلى لمنطقة التجميع.

ما الذي يفعله الحوت الكبير الذي يحتاج إلى تجميع صفقة شراء (Long) ضخمة أو تصريف كمياته؟

  • ضغط السعر: يرتفع السعر ببطء شديد دون إثارة الذعر بين أصحاب صفقات البيع، ولكنه يقترب خطوة بخطوة من منطقة التصفيات على الخريطة.
  • الاختراق الخاطف: بدفعة واحدة حادة، يتم سحب السعر لأعلى بنسبة 1.5% إلى 2% فوق مستوى المقاومة.
  • تفاعل تسلسلي: في هذه اللحظة، تتفعل أوامر الوقف الأولى، وتتطاير الصفقات ذات الرافعة العالية نحو التصفية. المنصة تشتري الأصل تلقائيًا بسعر السوق لإغلاق صفقات البيع هذه.
  • انفجار الحجم: طلبات الشراء بسعر السوق الناتجة عن تصفيات المنصة تصطدم بطلبات البيع المحددة (Limit Orders) التي وضعها الحوت مسبقًا لاستيعاب هذا الطلب الهستيري.

متداول التجزئة يعتقد: «اختراق! الاتجاه انعكس، سأدخل شراء فورًا!». يشتري في ذروة الزخم (FOMO)، ليبدأ الحوت الكبير في عكس موقفه تمامًا، ويهوي السعر كالحجر لضرب أوامر وقف صفقات الشراء هذه المرة.

لماذا ينفجر دفتر الطلبات (Order Book)

عندما تتفعل تصفية لصفقة حوت برافعة 20x بقيمة 500,000 دولار، تُصدر المنصة أمر سوق فوري لإغلاقها. إذا لم يكن هناك طلبات حدية (Limit Orders) معاكسة كافية عند تلك المستويات، ينساب السعر بحرية عبر المستويات.

  • ابتلاع طلبات الحد: محرك التصفية يبدأ بضرب دفتر الطلبات. إذا كان هناك 5 بيتكوين فقط عند مستوى 60,000 دولار بينما المطلوب تصفيته هو 20 بيتكوين، يخرق السعر المستوى ويستمر في الصعود حتى يجد سيولة كافية.
  • تأثير الدومينو: كل مستوى يتم اختراقه يفعل أوامر وقف وتصفيات جديدة لصفقات ذات رافعة أقل (مثل أصحاب رافعة 10x).
  • الارتداد العكسي (Dead Cat Bounce): بمجرد أن تحرق هذه الموجة كل السيولة المتاحة، يصبح دفتر الطلبات شبه فارغ. في هذه اللحظة، يدخل صانع السوق الذي كان ينتظر هذا الإنهاك بأوامر حدية ضخمة عند أفضل سعر، ليعكس اتجاه السعر فورًا.

نماذج التداول (Setups): كيف تسحب الأرباح من مناطق التصفية

لا تحاول التداول أثناء لحظة التصفية نفسها بشكل أعمى — السبريد الانزلاق السعري (Slippage) سيسحقان حسابك. العمل الحقيقي يبدأ بناءً على رد فعل السوق بعد أن تهدأ هذه الفوضى.

النموذج الأول: التداول على صيد السيولة (Stop Hunting)

  • الشروط: يظهر على خريطة التصفية جدار كثيف من صفقات الشراء أو البيع يبعد عن السعر الحالي بنسبة أقل من 1%.
  • الخطة: انتظر اختراق المستوى. لا تدخل الصفقة أثناء الاختراق نفسه.
  • التنفيذ: يندفع السعر بقوة خارج المستوى، ويزداد الحجم على إطار الدقيقة بشكل خيالي، بينما ينخفض الفائدة المفتوحة (OI) بشكل حاد (تم تصفية الجميع). بمجرد إغلاق الشمعة بفتيل طويل (Wick) عائدة إلى ما دون المستوى — افتح صفقة في الاتجاه المعاكس.
  • وقف الخسارة: يُوضع تمامًا خلف أقصى ذروة لفتيل الشمعة التي جمعت السيولة. إذا تم دفع السعر مرة أخرى بعد ذلك المستوى، فهذا يعني أنه كان اختراقًا حقيقيًا وليس اقتناصًا زائفًا للسيولة. يجب أن تكون مخاطرك محددة بصرامة.

أتمتة المتابعة: سكربت لرصد التحركات الشاذة

لكي لا تجلس أمام الشاشة طوال اليوم بانتظار الشلال، استخدم السكربت أدناه. يقوم بمراقبة تدفق التصفيات على المنصة في الوقت الفعلي ويطبع في الكونسول الأحجام غير الطبيعية للإغلاقات القسرية.

الرمز مكتوب بلغة Python الصافية دون مكتبات خارجية معقدة، يعمل بثبات ولا يتأثر بانقطاع الاتصال مع المنصة.

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."
        )

كيف تدمج هذا في سير عملك اليومي:

  • حدد الحد الأدنى (`threshold_usd`) بما يتناسب مع حجم رأس مالك والأصل الذي تتداوله. للعملات الكبيرة مثل ($BTC, $ETH)، فإن حداً يبدأ من 500,000 $ يكشف عن تفغيل أوامر الوقف الضخمة الحقيقية.
  • عند تفاعل السكربت، انتقل فورًا إلى رسم EXMON البياني وتحقق من التكتلات (Clusters). إذا حدثت تصفية ضخمة عند مستوى دعم أو مقاومة قوي — فهذه هي إشارتك الرئيسية للتحضير لصفقة عكس الاتجاه (Counter-trend).

لماذا تضلل خريطة التصفية المتداول الفردي وكيف تقرأها بذكاء مع الـ Funding Rates

أكبر خطأ يقع فيه المتداول المبتدئ هو النظر إلى خريطة التصفية (Liquidation Heatmap) بمعزل عن باقي المؤشرات. وجود تكتل كثيف من الخطوط الحمراء أو الخضراء على الرسم البياني لا يضمن لك حتمية الانعكاس. الخريطة تظهر لك الوقود المحتمل فقط، لكنها لا تخبرك بالاتجاه الذي سيسلكه المحرك الأساسي للسوق.

لكي تميز بين الاختراق الوهمي (Fakeout) واختراق السيولة الحقيقي (Liquidity Sweep) الممتد، عليك ربط معدل التمويل (Funding Rate) بالفائدة المفتوحة (Open Interest):

  • معدل تمويل موجب بشكل متطرف + ارتفاع الـ Open Interest: المتداولون يضخون صفقات شراء (Longs) بالروافع المالية بكثافة، والجميع مقتنع بالصعود. هذه هي البيئة المثالية لعملية "تصفية عنيفة" للحيتان لضرب صفقات الشراء والنزول بالسعر. خريطة التصفيات من الأسفل ستكون ممتلئة تماماً.
  • معدل تمويل سالب بشكل حاد + ارتفاع الـ Open Interest: السوق مشبع بصفقات البيع (Shorts). أي دفعة صعودية صغيرة قد تتسبب في شورت سكويز (Short Squeeze) متسلسل، لأن البائعين سيندفعون لإعادة شراء مراكزهم لإغلاقها.

إذا رأيت تكتلاً ضخماً لتصفيات الـ Shorts على الخريطة، وكان معدل التمويل سالباً بشدة لعدة ساعات متواصلة، فإن احتمالية حدوث صعود عنيف وتطهير للسيولة تتجاوز الحدود. الحيتان وكبار اللاعبين لن يتركوا هذه السيولة السهلة على الطاولة.

إعداد التداول #2: اختراق المستوى مع تدفق السيولة الحقيقي (استمرار الاتجاه)

ليس كل سحب للسيولة ينتهي بارتداد سريع وترك ذيل شمعة (Wick). في بعض الأحيان يكون حجم الوقود المتراكم ضخماً لدرجة أن السعر لا يكتفي بعمل ذيل تصفية، بل ينطلق في اتجاه وزخم قوي يستمر لعدة ساعات.

كيف تكتشف الاختراق الحقيقي لمنطقة التصفيات:

  • غياب الارتداد السريع: يخترق السعر مستوى تكتل أوامر الإيقاف، وبدلاً من ترك ذيل طويل، يغلق بشمعة ساعة كاملة وبجسم قوي فوق أو تحت المستوى.
  • سلوك مؤشر CVD (Cumulative Volume Delta): يتجه دلتا الشراء/البيع العدائي نحو ارتفاع عمودي صاروخي بالتوازي مع حركة السعر. هذا يعني أن من يشتري ليس فقط آليات التصفية الإجبارية للمنصة، بل إن مشتري السوق الكبار (Market Buyers) ينضمون بكل ثقلهم للحركة.
  • وضع الفائدة المفتوحة (OI) بعد الاختراق: بعد انتهاء موجة التصفيات، لا ينهار الـ OI إلى الصفر (أو ينخفض قليلاً ثم يعاود الارتفاع فوراً). هذه علامة جازمة على أن هناك مراكز ضخمة جديدة تدخل لتغطي المساحة المفرغة، وليس مجرد إغلاق لصفقات قديمة.

في هذه الحالة، تكون نقطة الدخول عند إعادة اختبار (Retest) المنطقة المخترقة، مع وضع أمر إيقاف خسارة (Stop Loss) صارم خلف حدود هذا المستوى.

قواعد صارمة لإدارة المخاطر أثناء اقتناص التصفيات

التداول داخل مناطق التصفيات يشبه السير في حقل ألغام. خطأ واحد في التوقيت بضع ثوانٍ أو الدخول بدون ستوب لوس، وستجد نفسك تحولت فوراً إلى مجرد نقطة بيانات على خريطة حرارية لشخص آخر.

  • انسَ الروافع المالية الأعلى من 5x–10x: عندما تتداول عكس الاندفاعات الفائقة التذبذب أثناء ضرب أوامر الإيقاف، قد يمنحك السوق الانزلاق السعري (Slippage) بنسبة 1-2% في أجزاء من الثانية. الرافعة العالية ستلتهم محفظتك قبل أن يتم تنفيذ أمر إيقاف الخسارة الخاص بك في دفتر الطلبات.
  • إيقاف الخسارة مقدس: إذا اخترق السعر مستوى التصفية وانطلق بعيداً عنك، لا تقم بالتبريد (Averaging down). تقبل خسارة نسبة محددة من رأس مالك (1-2%)، واعترف بأن السوق كان على حق، وانتظر منطقة إعداد جديدة ونظيفة.
  • تجاهل التكتلات الصغيرة: لا تتفاعل مع كل تقلب بسيط للسعر في دفتر الطلبات. صانع السوق والحيتان يصطادون الأموال الضخمة فقط. ركز على التكتلات التي يتجاوز حجمها متوسط حجم التداول لدقيقة واحدة لهذا الأصل بأضعاف مضاعفة.

التداول ليس مجرد تخمين لاتجاه الشمعة الخضراء القادمة. إنه حساب رياضياتي والاحتمالات واستغلال أخطاء الآخرين بدم بارد لصالحك. حلل خريطة السيولة، فكّر كصانع سوق، ولا تترك مركزك بدون حماية أبداً.

تلخيص هذه التدوينة باستخدام:

FAQ

خريطة التصفية هي أداة تحليل بياني تحسب مستويات الأسعار التقريبية التي تواجه عندها الصفقات الآجلة ذات الرافعة المالية التصفية الإجبارية. يقوم الخوارزمية بحساب مناطق تراكم أوامر السوق الإجبارية عن طريق دمج بيانات الفائدة المفتوحة (Open Interest) مع مستويات الرافعة المالية المختلفة (25x, 50x, 100x) لتحديد مناطق السيولة العالية.

تعتمد استراتيجية التداول على اقتناص السيولة (Liquidity Sweep) حيث تُشكل تكتلات التصفية الكبيرة أهدافاً رئيسية لحركة السعر. يضع المتداولون أوامر أخذ الأرباح داخل منطقة التصفية المتسلسلة مباشرةً، بينما يتم فتح صفقات الانعكاس فور استيعاب جميع الأوامر الإجبارية بواسطة الأوامر المحددة (Limit Orders).

توفر خرائط التصفية تقديرات احتمالية إحصائية وليست بيانات دقيقة مباشرة من سجلات التداول الخاصة بالمنصات. تقوم الأداة بنمذجة توزيع الرافعة المالية في السوق من خلال دمج بيانات الفائدة المفتوحة اللحظية، وتغيرات معدلات التمويل (Funding Rates)، وتدفقات أحجام الدلتا، مما يوفر دقة عالية في تحديد مستويات الدعم والمقاومة الهيكلية.
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

...

شاركنا برأيك

لن يتم نشر عنوان بريدك الإلكتروني. الحقول الإلزامية مشار إليها *