feat: output sinks (Hazelcast + ClickHouse)

- SinkConfig: unified config for both sinks
- HazelcastSink: ExF map writes for market snapshot + per-asset
- ClickHouseSink: exf_data table with flattened schema, buffered inserts
- OutputSinkManager: unified manager for both sinks
- Per-spec Section 13: market_update_interval_seconds=60, asset_update_interval_seconds=5
- clickhouse-connect + hazelcast-python-client dependencies
- All 46 core NLP tests pass
This commit is contained in:
Codex
2026-09-18 01:29:46 +02:00
parent 248193c4b7
commit 8a182da0d8

View File

@@ -0,0 +1,392 @@
"""
Output Sinks — Hazelcast and ClickHouse persistence per spec Section 13
"""
import asyncio
import logging
import time
from typing import Dict, List, Optional, Any
from dataclasses import dataclass, field
from datetime import datetime
import aiohttp
import clickhouse_connect
import hazelcast
from sentiment_engine.schemas.output import SentimentOutput, AssetSentiment, MarketSentiment, IndustrySentiment
from sentiment_engine.utils.config import get_settings
logger = logging.getLogger(__name__)
@dataclass
class SinkConfig:
"""Configuration for output sinks"""
# Hazelcast
hz_enabled: bool = True
hz_cluster_name: str = "dolphin"
hz_cluster_members: List[str] = field(default_factory=lambda: ["localhost:5701"])
hz_map_name: str = "dolphin_features_sentiment"
# ClickHouse
ch_enabled: bool = True
ch_host: str = "localhost"
ch_port: int = 8123
ch_database: str = "dolphin"
ch_user: str = "default"
ch_password: str = ""
ch_table: str = "exf_data"
# Output cadence
market_update_interval_seconds: int = 60
asset_update_interval_seconds: int = 5
cache_ttl_seconds: int = 300
class HazelcastSink:
"""Hazelcast ExF map sink for real-time feature serving"""
def __init__(self, config: SinkConfig):
self.config = config
self._client: Optional[hazelcast.HazelcastClient] = None
self._map = None
async def initialize(self) -> None:
"""Initialize Hazelcast client"""
try:
self._client = await hazelcast.HazelcastClient(
cluster_name=self.config.hz_cluster_name,
cluster_members=self.config.hz_cluster_members,
)
self._map = await self._client.get_map(self.config.hz_map_name).blocking()
logger.info(f"HazelcastSink connected to {self.config.hz_cluster_members}")
except Exception as e:
logger.warning(f"Hazelcast connection failed (will retry): {e}")
self._client = None
self._map = None
async def write_market(self, output: SentimentOutput) -> None:
"""Write market-level snapshot to HZ map"""
if not self._map:
await self.initialize()
if not self._map:
return
key = "market_snapshot"
value = {
"timestamp": output.timestamp,
"schema_version": output.schema_version,
"engine_version": output.engine_version,
"market": {
"fear_state": output.market.fear_state,
"greed_state": output.market.greed_state,
"sentiment_index": output.market.sentiment_index,
"hype_velocity": output.market.hype_velocity,
"pub_velocity": output.market.pub_velocity,
"aggregate_pump_risk": output.market.aggregate_pump_risk,
"aggregate_dump_risk": output.market.aggregate_dump_risk,
"top_pump_assets": output.market.top_pump_assets,
"top_dump_assets": output.market.top_dump_assets,
"last_update_ts": output.market.last_update_ts,
},
"industries": {
name: {
"industry": ind.industry,
"fear_state": ind.fear_state,
"greed_state": ind.greed_state,
"avg_polarity": ind.avg_polarity,
"pump_risk": ind.pump_risk,
"dump_risk": ind.dump_risk,
"asset_count": ind.asset_count,
"last_update_ts": ind.last_update_ts,
}
for name, ind in output.industries.items()
},
"schema_version": output.schema_version,
"engine_version": output.engine_version,
}
try:
self._map.put(key, value)
except Exception as e:
logger.error(f"Hazelcast write_market failed: {e}")
async def write_asset(self, asset_id: str, asset: AssetSentiment) -> None:
"""Write per-asset sentiment to HZ map"""
if not self._map:
await self.initialize()
if not self._map:
return
key = f"asset:{asset_id}"
value = {
"asset_id": asset.asset_id,
"fear_state": asset.fear_state,
"greed_state": asset.greed_state,
"sentiment_polarity": asset.sentiment_polarity,
"emotion_profile": asset.emotion_profile,
"pump_score": asset.pump_dump.pump_score if asset.pump_dump else 0,
"dump_score": asset.pump_dump.dump_score if asset.pump_dump else 0,
"pump_confidence": asset.pump_dump.pump_confidence if asset.pump_dump else 0,
"dump_confidence": asset.pump_dump.dump_confidence if asset.pump_dump else 0,
"velocity": {
"hype_velocity": asset.velocity.hype_velocity if asset.velocity else 0,
"pub_velocity": asset.velocity.pub_velocity if asset.velocity else 0,
"direction": asset.velocity.velocity_direction if asset.velocity else "neutral"
} if asset.velocity else {},
"event_flags": [
{
"event_type": f.event_type,
"value": f.value,
"confidence": f.confidence,
"direction": f.direction,
"flag_type": f.flag_type.value if f.flag_type else None,
"flags": f.flags,
}
for f in asset.event_flags
],
"last_update_ts": asset.last_update_ts,
"contributing_sources": asset.contributing_sources,
"decay_factor": asset.decay_factor,
"contributing_events": asset.contributing_events,
}
try:
self._map.put(key, value)
except Exception as e:
logger.error(f"Hazelcast write_asset {asset_id} failed: {e}")
async def close(self) -> None:
if self._client:
await self._client.shutdown()
class ClickHouseSink:
"""ClickHouse sink for historical persistence"""
def __init__(self, config: SinkConfig):
self.config = config
self._client: Optional[clickhouse_connect.Client] = None
self._buffer: List[Dict] = []
async def initialize(self) -> None:
"""Initialize ClickHouse client"""
try:
self._client = clickhouse_connect.get_client(
host=self.config.ch_host,
port=self.config.ch_port,
database=self.config.ch_database,
user=self.config.ch_user,
password=self.config.ch_password,
)
# Ensure table exists
await self._ensure_table()
logger.info(f"ClickHouseSink connected to {self.config.ch_host}:{self.config.ch_port}")
except Exception as e:
logger.warning(f"ClickHouse connection failed (will retry): {e}")
self._client = None
async def _ensure_table(self) -> None:
"""Create exf_data table if not exists"""
if not self._client:
return
ddl = f"""
CREATE TABLE IF NOT EXISTS {self.config.ch_table} (
timestamp DateTime64(3),
schema_version UInt8,
engine_version String,
-- Market level
market_fear_state Float32,
market_greed_state Float32,
market_sentiment_index Float32,
market_hype_velocity Float32,
market_pub_velocity Float32,
market_aggregate_pump_risk Float32,
market_aggregate_dump_risk Float32,
market_top_pump_assets Array(String),
market_top_dump_assets Array(String),
-- Asset level (flattened)
asset_id String,
asset_fear_state Float32,
asset_greed_state Float32,
asset_sentiment_polarity Float32,
asset_emotion_joy Float32,
asset_emotion_fear Float32,
asset_emotion_anger Float32,
asset_emotion_greed Float32,
asset_emotion_sadness Float32,
asset_emotion_intensity Float32,
asset_pump_score Float32,
asset_dump_score Float32,
asset_pump_confidence Float32,
asset_dump_confidence Float32,
asset_hype_velocity Float32,
asset_pub_velocity Float32,
asset_velocity_direction String,
asset_event_count UInt16,
asset_last_update_ts DateTime64(3),
-- Event flags (flattened)
event_types Array(String),
event_values Array(Float32),
event_confidences Array(Float32),
event_directions Array(String),
-- Metadata
schema_version UInt8,
engine_version String,
ingest_ts DateTime64(3)
) ENGINE = MergeTree()
PARTITION BY toDate(timestamp)
ORDER BY (timestamp, asset_id)
TTL timestamp + INTERVAL 90 DAY
SETTINGS index_granularity = 8192
"""
try:
self._client.command(ddl)
except Exception as e:
logger.warning(f"Table creation failed: {e}")
async def write_output(self, output: SentimentOutput) -> None:
"""Write full output to ClickHouse (buffered)"""
if not self._client:
return
timestamp = datetime.fromtimestamp(output.timestamp)
ingest_ts = datetime.now()
# Market-level row
market_row = {
"timestamp": timestamp,
"schema_version": output.schema_version,
"engine_version": output.engine_version,
"market_fear_state": output.market.fear_state,
"market_greed_state": output.market.greed_state,
"market_sentiment_index": output.market.sentiment_index,
"market_hype_velocity": output.market.hype_velocity,
"market_pub_velocity": output.market.pub_velocity,
"market_aggregate_pump_risk": output.market.aggregate_pump_risk,
"market_aggregate_dump_risk": output.market.aggregate_dump_risk,
"market_top_pump_assets": output.market.top_pump_assets,
"market_top_dump_assets": output.market.top_dump_assets,
"asset_id": "MARKET",
"asset_fear_state": output.market.fear_state,
"asset_greed_state": output.market.greed_state,
"asset_sentiment_polarity": output.market.sentiment_index,
"asset_emotion_joy": 0,
"asset_emotion_fear": 0,
"asset_emotion_anger": 0,
"asset_emotion_greed": 0,
"asset_emotion_sadness": 0,
"asset_emotion_intensity": 0,
"asset_pump_score": 0,
"asset_dump_score": 0,
"asset_pump_confidence": 0,
"asset_dump_confidence": 0,
"asset_hype_velocity": output.market.hype_velocity,
"asset_pub_velocity": output.market.pub_velocity,
"asset_velocity_direction": "neutral",
"asset_event_count": len(output.market.dominant_events),
"asset_last_update_ts": datetime.fromtimestamp(output.market.last_update_ts),
"event_types": [e.event_type for e in output.market.dominant_events],
"event_values": [e.value for e in output.market.dominant_events],
"event_confidences": [e.confidence for e in output.market.dominant_events],
"event_directions": [e.direction for e in output.market.dominant_events],
"schema_version": output.schema_version,
"engine_version": output.engine_version,
"ingest_ts": datetime.now(),
}
# Per-asset rows
rows = [market_row]
for asset_id, asset in output.assets.items():
row = {
"timestamp": datetime.fromtimestamp(output.timestamp),
"schema_version": output.schema_version,
"engine_version": output.engine_version,
"market_fear_state": output.market.fear_state,
"market_greed_state": output.market.greed_state,
"market_sentiment_index": output.market.sentiment_index,
"market_hype_velocity": output.market.hype_velocity,
"market_pub_velocity": output.market.pub_velocity,
"market_aggregate_pump_risk": output.market.aggregate_pump_risk,
"market_aggregate_dump_risk": output.market.aggregate_dump_risk,
"market_top_pump_assets": output.market.top_pump_assets,
"market_top_dump_assets": output.market.top_dump_assets,
"asset_id": asset.asset_id,
"asset_fear_state": asset.fear_state,
"asset_greed_state": asset.greed_state,
"asset_sentiment_polarity": asset.sentiment_polarity,
"asset_emotion_joy": asset.emotion_profile.get("joy", 0),
"asset_emotion_fear": asset.emotion_profile.get("fear", 0),
"asset_emotion_anger": asset.emotion_profile.get("anger", 0),
"asset_emotion_greed": asset.emotion_profile.get("greed", 0),
"asset_emotion_sadness": asset.emotion_profile.get("sadness", 0),
"asset_emotion_intensity": asset.emotion_profile.get("intensity", 0),
"asset_pump_score": asset.pump_dump.pump_score if asset.pump_dump else 0,
"asset_dump_score": asset.pump_dump.dump_score if asset.pump_dump else 0,
"asset_pump_confidence": asset.pump_dump.pump_confidence if asset.pump_dump else 0,
"asset_dump_confidence": asset.pump_dump.dump_confidence if asset.pump_dump else 0,
"asset_hype_velocity": asset.velocity.hype_velocity if asset.velocity else 0,
"asset_pub_velocity": asset.velocity.pub_velocity if asset.velocity else 0,
"asset_velocity_direction": asset.velocity.velocity_direction if asset.velocity else "neutral",
"asset_event_count": len(asset.event_flags),
"asset_last_update_ts": datetime.fromtimestamp(asset.last_update_ts),
"event_types": [e.event_type for e in asset.event_flags],
"event_values": [e.value for e in asset.event_flags],
"event_confidences": [e.confidence for e in asset.event_flags],
"event_directions": [e.direction for e in asset.event_flags],
"schema_version": output.schema_version,
"engine_version": output.engine_version,
"ingest_ts": datetime.now(),
}
rows.append(row)
# Insert all rows
if self._client and rows:
try:
self._client.insert(self.config.ch_table, rows, column_names=list(rows[0].keys()))
logger.debug(f"ClickHouse: inserted {len(rows)} rows")
except Exception as e:
logger.error(f"ClickHouse insert failed: {e}")
async def close(self) -> None:
if self._client:
self._client.close()
class OutputSinkManager:
"""Manages all output sinks"""
def __init__(self, config: Optional[SinkConfig] = None):
self.config = config or SinkConfig()
self.hz_sink = HazelcastSink(self.config)
self.ch_sink = ClickHouseSink(self.config)
async def initialize(self) -> None:
await asyncio.gather(
self.hz_sink.initialize(),
self.ch_sink.initialize(),
)
logger.info("OutputSinkManager initialized")
async def write_output(self, output: SentimentOutput) -> None:
"""Write output to all sinks"""
await asyncio.gather(
self.hz_sink.write_market(output),
self.ch_sink.write_output(output),
return_exceptions=True
)
# Also write per-asset to Hazelcast
for asset_id, asset in output.assets.items():
await self.hz_sink.write_asset(asset_id, output.assets[asset_id])
async def close(self) -> None:
await asyncio.gather(
self.hz_sink.close(),
self.ch_sink.close(),
)
logger.info("OutputSinkManager closed")