ニュース見出しをLLMでスコアリング|ドル円の急変を通知

AI×自動売買

先日、子供をお風呂に入れている20分ほどの間に、中東情勢の緊迫化を受けてドル円が急伸していました。上がったら乗ろうと思っていたのに、スマホを置いていた僕は完全に置いていかれました。「重要そうなニュースが出た瞬間だけ教えてくれる仕組みがあれば」と本気で思ったのが、今回作り始めたきっかけです。

僕は普段、ファンダメンタルズのニュースを四六時中追えるほど暇じゃありません(というか子育て中でそんな余裕はゼロです)。以前LINE通知やDiscord Webhookでテクニカル指標のアラートを組んだことはあったのですが、今回はニュース見出しそのものをLLMでスコアリングして、「本当にヤバそうな時だけ」通知を飛ばす仕組みを作りました。

前回は疑似コードで終わっていたので、今回はそのまま動く実装を書きます。ポイントは3つで、JSONを壊さないためのStructured Outputs、通知スパムを止める重複排除、そしてスコアが実際の値動きと合っているかの検証です。特に3つ目をやらないと、ただの高級な占いになります。

キーワード検索だけでは足りない理由

最初は「Iran」「FRB」「介入」みたいなキーワードを見出しに含むかどうかでフラグを立てようとしました。でもすぐに限界に気づきます。「停戦合意」というキーワードは本来ポジティブですが、「停戦崩れる懸念」だと真逆の意味になる。単語の有無だけを見ていると、こうした文脈の反転を拾えません。

もう1つの問題が量です。素朴なキーワードマッチだと、1日に数十件フラグが立ちます。全部に通知が飛ぶと、3日でミュートします。アラートは、鳴りすぎた瞬間に価値がゼロになるのが厄介なところです。

見出しを取得する

import hashlib
import re
import feedparser

RSS_URLS = [
    "https://news.google.com/rss/search?q=%E3%83%89%E3%83%AB%E5%86%86&hl=ja&gl=JP&ceid=JP:ja",
    "https://news.google.com/rss/search?q=USDJPY+OR+FRB&hl=en-US&gl=US&ceid=US:en",
]

def normalize(title: str) -> str:
    t = re.sub(r"\s+", "", title)
    t = re.sub(r"[||\-–—].*$", "", t)   # 媒体名の後ろを落とす
    return t.lower()

def fetch_headlines(limit: int = 15) -> list[dict]:
    seen, out = set(), []
    for url in RSS_URLS:
        feed = feedparser.parse(url)
        for entry in feed.entries[:limit]:
            key = hashlib.sha1(normalize(entry.title).encode()).hexdigest()[:16]
            if key in seen:
                continue
            seen.add(key)
            out.append({
                "id": key,
                "title": entry.title,
                "link": getattr(entry, "link", ""),
                "published": getattr(entry, "published", ""),
            })
    return out

同じニュースが複数の媒体から配信されるので、normalizeで表記ゆれを潰してからハッシュを取ります。ここを入れないと、同じ内容で3回通知が来ます。実際、最初のバージョンはそれで一晩に11回鳴りました。

Structured Outputsでスコアを受け取る

LLMに「JSONで返して」と頼むだけだと、たまに前置きの文章が付いてきてjson.loadsが落ちます。OpenAIのStructured Outputsを使うと、スキーマ準拠が保証されます。

import json
import os
from openai import OpenAI

client = OpenAI(api_key=os.getenv("OPENAI_API_KEY"))

SCHEMA = {
    "type": "object",
    "properties": {
        "items": {
            "type": "array",
            "items": {
                "type": "object",
                "properties": {
                    "id": {"type": "string"},
                    "score": {"type": "number"},      # -1.0=円高要因 / +1.0=円安要因
                    "confidence": {"type": "number"}, # 0.0-1.0
                    "reason": {"type": "string"},
                },
                "required": ["id", "score", "confidence", "reason"],
                "additionalProperties": False,
            },
        }
    },
    "required": ["items"],
    "additionalProperties": False,
}

SYSTEM = (
    "あなたは為替のニュース評価器です。各見出しがUSD/JPYに与える方向と強さを "
    "-1.0(強い円高要因)〜 +1.0(強い円安要因)で評価します。"
    "相場と無関係な見出しは score=0.0, confidence=0.0 としてください。"
    "推測が困難な場合は confidence を低くし、無理にスコアを付けないこと。"
)

def score_batch(headlines: list[dict]) -> dict[str, dict]:
    if not headlines:
        return {}
    payload = "\n".join(f'{h["id"]}\t{h["title"]}' for h in headlines)
    res = client.chat.completions.create(
        model="gpt-4o-mini",
        temperature=0,
        messages=[
            {"role": "system", "content": SYSTEM},
            {"role": "user", "content": f"id\tタイトル の形式です。全件評価してください。\n{payload}"},
        ],
        response_format={
            "type": "json_schema",
            "json_schema": {"name": "scores", "strict": True, "schema": SCHEMA},
        },
    )
    data = json.loads(res.choices[0].message.content)
    return {item["id"]: item for item in data["items"]}

1件ずつではなく、まとめて1リクエストで投げています。15件を個別に呼ぶと15回分の待ち時間とオーバーヘッドがかかりますが、まとめれば1回です。コストも待ち時間も1桁減りました。

temperature=0も必須です。同じ見出しに毎回違うスコアが付くと、閾値の意味がなくなります。

鳴りすぎを止める

閾値を超えただけで通知すると、大きなニュースの日に連発します。既読管理とクールダウンを入れます。

import sqlite3
import time
from pathlib import Path

DB = Path("news_alert.db")

def init_db() -> sqlite3.Connection:
    con = sqlite3.connect(DB)
    con.execute("""
        CREATE TABLE IF NOT EXISTS seen (
            id TEXT PRIMARY KEY,
            title TEXT,
            score REAL,
            confidence REAL,
            reason TEXT,
            ts INTEGER,
            notified INTEGER DEFAULT 0
        )
    """)
    con.commit()
    return con

COOLDOWN_SEC = 30 * 60

def should_notify(con: sqlite3.Connection, item: dict,
                  threshold: float = 0.6, min_conf: float = 0.5) -> bool:
    if abs(item["score"]) < threshold or item["confidence"] < min_conf:
        return False
    cur = con.execute("SELECT 1 FROM seen WHERE id = ?", (item["id"],))
    if cur.fetchone():
        return False   # 既読
    last = con.execute(
        "SELECT MAX(ts) FROM seen WHERE notified = 1 AND score * ? > 0",
        (item["score"],),
    ).fetchone()[0]
    if last and time.time() - last < COOLDOWN_SEC:
        return False   # 同方向は30分に1回まで
    return True

score * ? > 0で同じ符号のアラートだけをクールダウン対象にしています。円安方向の通知を出した直後に、逆方向の重要ニュースが来たら、それは鳴らすべきだからです。

通知して記録する

import requests

WEBHOOK = os.getenv("DISCORD_WEBHOOK_URL")

def send_discord(text: str) -> None:
    if not WEBHOOK:
        print("[dry-run]", text)
        return
    r = requests.post(WEBHOOK, json={"content": text[:1900]}, timeout=15)
    r.raise_for_status()

def run() -> None:
    con = init_db()
    headlines = fetch_headlines()
    fresh = [h for h in headlines
             if not con.execute("SELECT 1 FROM seen WHERE id=?", (h["id"],)).fetchone()]
    scores = score_batch(fresh)

    for h in fresh:
        s = scores.get(h["id"])
        if not s:
            continue
        notify = should_notify(con, {**s, "id": h["id"]})
        if notify:
            arrow = "円安↑" if s["score"] > 0 else "円高↓"
            send_discord(
                f"[{arrow} {s['score']:+.2f} / 確度{s['confidence']:.2f}]\n"
                f"{h['title']}\n{s['reason']}\n{h['link']}"
            )
        con.execute(
            "INSERT OR REPLACE INTO seen VALUES (?,?,?,?,?,?,?)",
            (h["id"], h["title"], s["score"], s["confidence"],
             s["reason"], int(time.time()), int(notify)),
        )
    con.commit()

if __name__ == "__main__":
    run()

通知しなかったものもDBに残しているのが大事です。これが次の検証の材料になります。閾値を後から下げたくなったとき、記録がなければ「下げたらどうなっていたか」を確かめられません。

スコアが当たっているか検証する

ここが一番やるべきことでした。1〜2週間ためたら、スコアとその後のドル円の動きを突き合わせます。

import pandas as pd
import yfinance as yf

def validate(con: sqlite3.Connection, horizon_min: int = 60) -> pd.DataFrame:
    rows = pd.read_sql("SELECT * FROM seen WHERE confidence >= 0.5", con)
    if rows.empty:
        return rows
    rows["dt"] = pd.to_datetime(rows["ts"], unit="s")

    fx = yf.download("JPY=X", period="1mo", interval="15m",
                     auto_adjust=False, progress=False)
    if isinstance(fx.columns, pd.MultiIndex):
        fx.columns = fx.columns.get_level_values(0)
    px = fx["Close"].dropna()
    px.index = pd.to_datetime(px.index).tz_localize(None)

    out = []
    for _, r in rows.iterrows():
        after = px[px.index >= r["dt"]]
        if len(after) < 2:
            continue
        base = float(after.iloc[0])
        later = after[after.index <= r["dt"] + pd.Timedelta(minutes=horizon_min)]
        if len(later) < 2:
            continue
        move = float(later.iloc[-1]) / base - 1
        out.append({
            "score": r["score"],
            "move_pct": move * 100,
            "hit": (move > 0) == (r["score"] > 0),
        })
    df = pd.DataFrame(out)
    if df.empty:
        return df
    strong = df[df["score"].abs() >= 0.6]
    print(f"全体 n={len(df)} 方向一致率={df['hit'].mean():.1%}")
    print(f"強シグナルのみ n={len(strong)} 一致率={strong['hit'].mean():.1%}")
    print(f"スコアと変動の相関={df['score'].corr(df['move_pct']):.3f}")
    return df

僕の手元では、方向一致率は5割台でした。つまりスコアの符号で方向を当てるのはほぼ無理です。ニュース見出しの感情スコアが翌日リターンをうまく説明しないのは、株式ニュースの感情分析を検証した記事でも同じ結論でした。

ただし、move_pct絶対値を見ると話が変わります。強シグナルが出た直後の1時間は、平常時より値幅が出ていました。つまりこの仕組みが教えてくれるのは「どっちに動くか」ではなく「今、動いている」です。それで十分でした。お風呂から出てチャートを開くきっかけになれば目的は果たしています。

コストと運用

項目設定ねらい
実行間隔15分ごと(cron / タスクスケジューラ)秒単位の判断はそもそも狙わない
1回の件数新規のみ、多くて15件既読はAPIに投げない
モデル小さいモデル + temperature=0分類タスクに大型は不要
閾値score 0.6 かつ confidence 0.5鳴りすぎ防止
クールダウン同方向30分大ニュース時の連発を抑える

既読分をAPIに投げない設計にしたので、1日の呼び出しは数十件ぶんに収まりました。APIキーはos.getenvで読み、.env.gitignoreに入れています。うっかりコミットすると、キーは即座に無効化される羽目になります。

まとめ

ニュース見出しのスコアリングは、キーワード検索の「文脈を読めない」弱点をかなり補ってくれました。ただ、実際に作ってみて分かったのは、難しいのはスコアリングそのものではなく、鳴りすぎを止めることと、スコアを信じてよいか確かめることだという点です。

やったことは4つです。バッチで1リクエストにまとめる。Structured Outputsで壊れないJSONを受け取る。ハッシュ重複排除とクールダウンで通知を絞る。そして貯めたスコアを実際の値動きと突き合わせる。

方向は当たりませんでしたが、「動いている」ことは分かります。それだけでも、お風呂上がりにチャートを開く理由にはなります。次はドル円だけでなく、保有している日本株の銘柄名でも同じ仕組みを回してみるつもりです。

タイトルとURLをコピーしました