peyterho's picture
Upload macro_sentiment/pipeline.py
cbdd80d verified
Raw
History Blame
6.78 kB
"""
Unified Macroeconomic Sentiment Pipeline.
Combines dictionary signals + transformer ensemble + topic routing
into a structured output with overall macro sentiment, policy stance,
crisis signal, and domain-specific scores.
"""
import json
import numpy as np
from typing import Dict, List
from dataclasses import dataclass, field, asdict
from macro_sentiment.dictionaries import CombinedDictionaryScorer
from macro_sentiment.transformers_ensemble import TransformerEnsemble
@dataclass
class MacroSentimentResult:
"""Structured output from the macro sentiment pipeline."""
macro_sentiment: float = 0.0
confidence: float = 0.0
financial_sentiment: float = 0.0
policy_stance: float = 0.0
climate_sentiment: float = 0.0
crisis_signal: float = 0.0
uncertainty: float = 0.0
detected_domain: str = "general"
topic: str = ""
topic_confidence: float = 0.0
head_used: str = ""
lm_polarity: float = 0.0
henry_polarity: float = 0.0
climate_exposure: float = 0.0
raw_features: Dict[str, float] = field(default_factory=dict)
def to_dict(self): return asdict(self)
def to_json(self, indent=2): return json.dumps(self.to_dict(), indent=indent, default=str)
def summary(self):
s = "Very Negative" if self.macro_sentiment < -0.6 else "Negative" if self.macro_sentiment < -0.2 else "Neutral" if self.macro_sentiment < 0.2 else "Positive" if self.macro_sentiment < 0.6 else "Very Positive"
p = "Very Dovish" if self.policy_stance < -0.6 else "Dovish" if self.policy_stance < -0.2 else "Neutral" if self.policy_stance < 0.2 else "Hawkish" if self.policy_stance < 0.6 else "Very Hawkish"
c = "HIGH CRISIS" if self.crisis_signal > 0.6 else "Elevated" if self.crisis_signal > 0.3 else "Normal"
parts = [f"Sentiment: {s} ({self.macro_sentiment:+.3f})", f"Policy: {p} ({self.policy_stance:+.3f})", f"Crisis: {c} ({self.crisis_signal:.3f})", f"Domain: {self.detected_domain}"]
if self.climate_exposure > 0.1:
cl = "Opportunity" if self.climate_sentiment > 0.2 else "Risk" if self.climate_sentiment < -0.2 else "Neutral"
parts.append(f"Climate: {cl} (exp={self.climate_exposure:.2f})")
return " | ".join(parts)
class MacroSentimentPipeline:
def __init__(self, device="cpu", use_router=True, load_transformers=True):
self.device = device
self.use_transformers = load_transformers
self.dict_scorer = CombinedDictionaryScorer()
self.transformer_ensemble = TransformerEnsemble(device=device, use_router=use_router) if load_transformers else None
def __call__(self, text, mode="routed"):
return self.score(text, mode=mode)
def score(self, text, mode="routed"):
result = MacroSentimentResult()
dict_features = self.dict_scorer.score(text)
result.lm_polarity = dict_features["lm_polarity"]
result.henry_polarity = dict_features["henry_polarity"]
result.climate_exposure = dict_features["climate_exposure"]
result.crisis_signal = dict_features["macro_crisis_intensity"]
result.uncertainty = dict_features["macro_uncertainty"]
if self.use_transformers and mode != "dict_only":
if mode == "routed":
tf_features = self.transformer_ensemble.score_routed(text)
result.head_used = tf_features.get("head_used", "unknown")
result.topic = tf_features.get("topic", "")
result.topic_confidence = tf_features.get("topic_confidence", 0.0)
head = result.head_used
result.detected_domain = {"policy": "policy", "climate": "climate", "tweet": "social"}.get(head, "financial_news")
primary_tf_score = tf_features.get("sentiment_score", 0.0)
elif mode == "all":
tf_features = self.transformer_ensemble.score_all(text)
primary_tf_score = tf_features.get("ensemble_mean", 0.0)
result.detected_domain = "ensemble"
result.head_used = "all"
else:
tf_features = {}
primary_tf_score = 0.0
result.financial_sentiment = tf_features.get("finbert_composite_score", tf_features.get("sentiment_score", 0.0))
result.climate_sentiment = tf_features.get("climate_composite_score", dict_features.get("climate_net_sentiment", 0.0))
policy_score = tf_features.get("policy_composite_score", None)
dict_policy = dict_features["macro_policy_stance"]
if dict_policy != 0.0:
result.policy_stance = 0.8 * dict_policy + 0.2 * (policy_score or 0.0)
elif policy_score is not None:
result.policy_stance = policy_score
crisis_weight = min(0.4, result.crisis_signal * 0.5)
tf_weight = 0.65 - crisis_weight
dict_weight = 0.35 + crisis_weight
dict_composite = np.mean([dict_features["lm_polarity"], dict_features["henry_polarity"]])
result.macro_sentiment = tf_weight * primary_tf_score + dict_weight * dict_composite
signals = [primary_tf_score, dict_composite]
agreement = 1.0 - np.std(signals)
result.confidence = min(1.0, max(0.0, 0.5 * agreement + 0.3 * result.topic_confidence + 0.2 * (1.0 - result.uncertainty)))
else:
result.detected_domain = "dict_only"
result.head_used = "none"
dict_composite = np.mean([dict_features["lm_polarity"], dict_features["henry_polarity"]])
result.macro_sentiment = dict_composite
result.policy_stance = dict_features["macro_policy_stance"]
result.climate_sentiment = dict_features["climate_net_sentiment"]
result.financial_sentiment = dict_composite
result.confidence = 0.3
for attr in ["macro_sentiment", "financial_sentiment", "policy_stance", "climate_sentiment"]:
setattr(result, attr, float(np.clip(getattr(result, attr), -1.0, 1.0)))
for attr in ["crisis_signal", "uncertainty", "confidence"]:
setattr(result, attr, float(np.clip(getattr(result, attr), 0.0, 1.0)))
result.raw_features = {**dict_features}
if self.use_transformers and mode != "dict_only":
result.raw_features.update({f"tf_{k}": v for k, v in tf_features.items() if isinstance(v, (int, float))})
return result
def score_batch(self, texts, mode="routed"):
return [self.score(t, mode=mode) for t in texts]
def load_all(self):
if self.transformer_ensemble:
self.transformer_ensemble.load_all()
return self