1
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

タイのデイリーニュースを取得しタイ...EventBridge, Lambda, Bedrockによるイベント駆動型アーキテクチャ

1
Posted at

はじめに

最近、海外(タイ)に関わる情報を追う必要が出てきました。
現地の政治・経済・テックの動きは、案件の意思決定やコミュニケーションにも影響することがあるため、最低限のトレンドは押さえておきたいところです。

とはいえ、毎日自分でニュースを探して取捨選択して…という運用は、忙しい時ほど続きにくいのが現実です。
そこで「収集と要約はAIに任せて、読む判断だけ自分がやる」形に寄せたくなり、この仕組みを作りました。

アーキテクチャ構成

image.png

  • 日本時間の9時に毎日動くようにEventBridgeで設定します
  • EventBridgeは、時間になったらLambdaを起動します
  • Lambdaは下記の処理を実施します
    • RSSやAPIによってタイの情報(政治、経済、テック)を読み取る
    • Bedrockに要約を依頼
    • 要約後、マークダウン形式でS3に保存
  • ユーザはLINEで通知を受け取り、Streamlitの画面から今日のニュースを確認できます ※実装前

実装

Lambda関数

関数名とランタイムの設定を行います。

環境変数 値(参考値)
BEDROCK_MODEL_ID us.anthropic.claude-3-7-sonnet-20250219-v1:0
BEDROCK_REGION us-west-2
FETCH_TIMEOUT_SEC 25
HTTP_USER_AGENT global-news-bot/0.1
MAX_ITEMS_PER_CATEGORY 30
MAX_ITEMS_PER_FEED 20
S3_BUCKET s3-bucket-sample

image.png

lambda_handler.pyは下記の実装とします。
boto3とfeedparser、requestsについてはzip化して一緒にアップロードしました。

lambda_function.py
# lambda_function.py
import os
import json
import re
import html
import hashlib
from datetime import datetime
from zoneinfo import ZoneInfo
from typing import List, Dict, Any, Tuple, Optional
from urllib.parse import urlparse, urlunparse, parse_qsl, urlencode

import boto3
import feedparser
import requests


# ---------- Config ----------
APP_PREFIX = "global-news-"

S3_BUCKET = os.environ["S3_BUCKET"]

MAX_ITEMS_PER_CATEGORY = int(os.environ.get("MAX_ITEMS_PER_CATEGORY", "30"))
MAX_ITEMS_PER_FEED = int(os.environ.get("MAX_ITEMS_PER_FEED", "20"))
FETCH_TIMEOUT_SEC = int(os.environ.get("FETCH_TIMEOUT_SEC", "10"))
HTTP_USER_AGENT = os.environ.get("HTTP_USER_AGENT", f"{APP_PREFIX}bot/0.1")

BEDROCK_REGION = os.environ.get("BEDROCK_REGION", os.environ.get("AWS_REGION", "us-west-2"))
BEDROCK_MODEL_ID = os.environ.get("BEDROCK_MODEL_ID", "amazon.nova-pro-v1:0")

# RSS feeds (comma-separated). If not provided, default to the curated set below.
from datetime import datetime
from zoneinfo import ZoneInfo


TH_TZ = ZoneInfo("Asia/Bangkok")


def gdelt_timerange():
    today = datetime.now(TH_TZ)

    start = today.replace(hour=0, minute=0, second=0)
    end   = today.replace(hour=23, minute=59, second=59)

    return (
        start.strftime("%Y%m%d%H%M%S"),
        end.strftime("%Y%m%d%H%M%S"),
    )


def build_gdelt_url(query: str) -> str:
    start, end = gdelt_timerange()

    return (
        "https://api.gdeltproject.org/api/v2/doc/doc"
        f"?query={query}"
        "&mode=ArtList"
        "&format=rss"
        "&maxrecords=50"
        "&sort=HybridRel"
        f"&startdatetime={start}"
        f"&enddatetime={end}"
    )


# ============================
# RSS 定義(1行ずつ)
# ============================

DEFAULT_RSS_POLITICS = [
    build_gdelt_url("thailand (politics OR government OR election)")
]

DEFAULT_RSS_ECONOMY = [
    build_gdelt_url("thailand (economy OR gdp OR inflation OR bank OR baht OR trade OR tourism)")
]

DEFAULT_RSS_TECH = [
    build_gdelt_url("thailand (technology OR ai OR cyber OR software OR startup OR digital)")
]

# Date boundary: Thailand time is natural for Thailand news daily cut
TH_TZ = ZoneInfo("Asia/Bangkok")


# ---------- Clients ----------
s3 = boto3.client("s3")
brt = boto3.client("bedrock-runtime", region_name=BEDROCK_REGION)


# ---------- Helpers ----------
def _clean_text(s: str) -> str:
    if not s:
        return ""
    s = html.unescape(s)
    s = re.sub(r"<[^>]+>", " ", s)  # strip HTML tags
    s = re.sub(r"\s+", " ", s).strip()
    return s


def _hash(s: str) -> str:
    return hashlib.sha256(s.encode("utf-8")).hexdigest()[:16]


def _normalize_url(url: str) -> str:
    """
    Normalize URL to improve dedup:
    - strip common tracking query params (utm_*, fbclid, etc.)
    - keep stable parts
    """
    if not url:
        return url
    try:
        p = urlparse(url)
        q = []
        for k, v in parse_qsl(p.query, keep_blank_values=True):
            lk = k.lower()
            if lk.startswith("utm_"):
                continue
            if lk in ("fbclid", "gclid", "igshid", "mc_cid", "mc_eid"):
                continue
            q.append((k, v))
        new_query = urlencode(q, doseq=True)
        return urlunparse((p.scheme, p.netloc, p.path, p.params, new_query, p.fragment))
    except Exception:
        return url


def _fetch_url(url: str) -> Optional[bytes]:
    try:
        r = requests.get(
            url,
            headers={
                "User-Agent": HTTP_USER_AGENT,
                "Accept": "application/rss+xml, application/atom+xml, application/xml, text/xml, */*",
                "Accept-Encoding": "gzip, deflate",
                "Connection": "close",
            },
            timeout=(3, FETCH_TIMEOUT_SEC),
            allow_redirects=True,
        )

        ct = (r.headers.get("Content-Type") or "").lower()
        b = r.content or b""
        head = b[:200].decode("utf-8", errors="replace").replace("\n", " ").replace("\r", " ")

        print(
            f"[INFO] fetch url={url} status={r.status_code} bytes={len(b)} ct={ct} final_url={r.url}"
        )

        # 200だけどHTMLが返ってる(= botブロック/エラーページ)を検知
        if "text/html" in ct or head.lstrip().startswith("<!doctype html") or head.lstrip().startswith("<html"):
            print(f"[WARN] non-rss response (looks like HTML). url={url} head={head[:200]}")
            return None

        if 200 <= r.status_code < 300 and b:
            return b

        print(f"[WARN] fetch failed status={r.status_code} url={url} head={head[:200]}")
        return None

    except Exception as ex:
        print(f"[WARN] fetch exception url={url} ex={ex}")
        return None

def fetch_rss_items(feed_urls: List[str], max_items_per_feed: int, max_items_total: int) -> List[Dict[str, Any]]:
    """
    - Parse RSS/Atom robustly
    - Dedup by normalized URL
    - Return up to max_items_total
    """
    items: List[Dict[str, Any]] = []

    for url in feed_urls:
        raw = _fetch_url(url)
        if not raw:
            continue

        try:
            fp = feedparser.parse(raw)
            entries = getattr(fp, "entries", []) or []

            print(f"[INFO] parsed feed url={url} bozo={getattr(fp, 'bozo', None)} entries={len(entries)}")
            if getattr(fp, "bozo", 0):
                print(f"[WARN] feed bozo url={url} ex={getattr(fp, 'bozo_exception', None)}")

            for e in entries[: max_items_per_feed]:
                title = _clean_text(getattr(e, "title", ""))
                link = getattr(e, "link", "") or ""
                link = _normalize_url(link)
                summary = _clean_text(getattr(e, "summary", "") or getattr(e, "description", ""))
                published = getattr(e, "published", "") or getattr(e, "updated", "") or ""

                if not title or not link:
                    continue

                items.append(
                    {
                        "source_feed": url,
                        "title": title,
                        "link": link,
                        "summary": summary,
                        "published": published,
                        "id": _hash(link),
                    }
                )
        except Exception as ex:
            print(f"[WARN] RSS parse failed url={url} ex={ex}")

    # Dedup by normalized link
    seen = set()
    dedup = []
    for it in items:
        if it["link"] in seen:
            continue
        seen.add(it["link"])
        dedup.append(it)

    # Prefer items that have some summary
    dedup.sort(key=lambda x: (0 if x.get("summary") else 1, x.get("published", "")), reverse=False)

    return dedup[:max_items_total]


def bedrock_summarize_and_translate(category_name_ja: str, items: List[Dict[str, Any]]) -> str:
    """
    PoC: summarize based on title/summary only.
    Assumes Nova-like Bedrock message schema.
    """
    compact = []
    for i, it in enumerate(items, start=1):
        compact.append(
            {
                "no": i,
                "title": it["title"],
                "summary": (it["summary"] or "")[:800],
                "published": it["published"],
                "url": it["link"],
                "source_feed": it["source_feed"],
            }
        )

    prompt = f"""
あなたは国際ニュース編集者です。以下はタイのニュース(カテゴリ: {category_name_ja})のRSS抜粋です。

要件:
- 出力は日本語
- Markdown形式
- 最初に「今日の要点(3〜6点)」を箇条書き
- 次に「記事一覧」として、各記事を見出し付きで要約(2〜4行)し、最後に必ずURLを記載
- 不確かな推測は禁止。与えられた情報(title/summary)から言える範囲で書く
- 誇張せず、事実ベースで簡潔に
- 同じ話題が複数記事にある場合は、記事一覧は残しつつ「同一トピック」と分かるように表現を揃える

入力データ(JSON):
{json.dumps(compact, ensure_ascii=False)}
""".strip()

    body = {
        "anthropic_version": "bedrock-2023-05-31",
        "max_tokens": 1800,
        "temperature": 0.2,
        "messages": [
            {
                "role": "user",
                "content": [
                    {"type": "text", "text": prompt}
                ]
            }
        ]
    }

    resp = brt.invoke_model(
        modelId=BEDROCK_MODEL_ID,
        body=json.dumps(body).encode("utf-8"),
        accept="application/json",
        contentType="application/json",
    )

    payload = json.loads(resp["body"].read().decode("utf-8"))

    # ✅ Claude(Anthropic)の返却は payload["content"] が block配列になりがち
    # 例: [{"type":"text","text":"...markdown..."}, ...]
    out = ""

    # 1) Claude形式(content blocks)
    content_blocks = payload.get("content")
    if isinstance(content_blocks, list):
        parts = []
        for b in content_blocks:
            if isinstance(b, dict) and b.get("type") == "text":
                t = b.get("text", "")
                if isinstance(t, str) and t:
                    parts.append(t)
        out = "\n".join(parts).strip()

    # 2) Converse/Nova系フォールバック(output.message.content)
    if not out:
        parts = []
        for b in payload.get("output", {}).get("message", {}).get("content", []):
            if isinstance(b, dict) and "text" in b and isinstance(b["text"], str):
                parts.append(b["text"])
        out = "\n".join(parts).strip()

    # 3) 最終フォールバック(completion等)
    if not out:
        if isinstance(payload.get("completion"), str):
            out = payload["completion"].strip()
        elif isinstance(payload.get("output_text"), str):
            out = payload["output_text"].strip()

    # 4) それでも空なら、payloadをダンプ(デバッグ用)※本番では消してOK
    if not out:
        out = json.dumps(payload, ensure_ascii=False)

    return out



def build_daily_markdown(date_str: str, sections: List[Tuple[str, str]]) -> str:
    header = f"# Thailand Daily News ({date_str})\n\n"
    toc = "## 目次\n" + "\n".join([f"- [{title}](#{title})" for title, _ in sections]) + "\n\n"
    body = ""
    for title, md in sections:
        body += f"## {title}\n\n{md}\n\n---\n\n"
    return header + toc + body


def put_to_s3(markdown: str, date_str: str) -> str:
    # Requirement: Thailand/yyyy_mm_dd.md
    key = f"Thailand/{date_str}.md"
    s3.put_object(
        Bucket=S3_BUCKET,
        Key=key,
        Body=markdown.encode("utf-8"),
        ContentType="text/markdown; charset=utf-8",
        CacheControl="no-cache",
    )
    return key


# ---------- Lambda handler ----------
def lambda_handler(event, context):
    print(f"[INFO] timeout={FETCH_TIMEOUT_SEC} max_items_per_feed={MAX_ITEMS_PER_FEED} max_items_per_category={MAX_ITEMS_PER_CATEGORY}")

    today = datetime.now(TH_TZ)
    date_str = today.strftime("%Y_%m_%d")

    print(f"[INFO] start date={date_str} bucket={S3_BUCKET} model={BEDROCK_MODEL_ID}")
    print(f"[INFO] feeds politics={len(RSS_POLITICS)} economy={len(RSS_ECONOMY)} tech={len(RSS_TECH)}")

    # 1) Fetch
    politics_items = fetch_rss_items(RSS_POLITICS, MAX_ITEMS_PER_FEED, MAX_ITEMS_PER_CATEGORY)
    economy_items = fetch_rss_items(RSS_ECONOMY, MAX_ITEMS_PER_FEED, MAX_ITEMS_PER_CATEGORY)
    tech_items = fetch_rss_items(RSS_TECH, MAX_ITEMS_PER_FEED, MAX_ITEMS_PER_CATEGORY)

    # 2) Summarize/Translate per category
    sections: List[Tuple[str, str]] = []

    if politics_items:
        sections.append(("政治", bedrock_summarize_and_translate("政治", politics_items)))
    else:
        sections.append(("政治", "_(取得0件:RSSが落ちている/フィード形式変更の可能性)_"))

    if economy_items:
        sections.append(("経済", bedrock_summarize_and_translate("経済", economy_items)))
    else:
        sections.append(("経済", "_(取得0件:RSSが落ちている/検索条件が強すぎる可能性)_"))

    if tech_items:
        sections.append(("テック", bedrock_summarize_and_translate("テック", tech_items)))
    else:
        sections.append(("テック", "_(取得0件:RSSが落ちている/フィード形式変更の可能性)_"))

    # 3) Build Markdown & Save to S3
    md = build_daily_markdown(date_str, sections)
    key = put_to_s3(md, date_str)

    return {
        "statusCode": 200,
        "body": json.dumps(
            {
                "message": "ok",
                "date": date_str,
                "s3_bucket": S3_BUCKET,
                "s3_key": key,
                "counts": {
                    "politics": len(politics_items),
                    "economy": len(economy_items),
                    "tech": len(tech_items),
                },
                "feeds": {
                    "politics": RSS_POLITICS,
                    "economy": RSS_ECONOMY,
                    "tech": RSS_TECH,
                },
            },
            ensure_ascii=False,
        ),
    }


S3バケット

バケット名を指定して作成します。
image.png

コストを考え、30日間でファイルが削除されるようにライフサイクルルールを設定します。
image.png

EventBridge

cronで毎日9時に実行されるように設定します。
image.png

Lambda関数を実行するように設定します。
image.png

フロントエンド

app.py
from dataclasses import dataclass
from datetime import date, datetime
from zoneinfo import ZoneInfo

import boto3, os, calendar
import streamlit as st
from botocore.exceptions import ClientError
from dotenv import load_dotenv


# ---------- Constants ----------
APP_PREFIX = "global-news-"
TH_TZ = ZoneInfo("Asia/Bangkok")  # Daily cut aligns with Thailand time


# ---------- Env Loader ----------
# ローカル実行時のみ .env を読む(Cloud では通常 .env は無い)
load_dotenv(override=False)


def get_env(key: str, default: str | None = None) -> str | None:
    """
    Priority:
      1. OS env (.env 포함)
      2. st.secrets (Streamlit Cloud)
      3. default
    """
    if key in os.environ:
        return os.environ.get(key)

    if hasattr(st, "secrets") and key in st.secrets:
        return st.secrets.get(key)

    return default


# ---------- Config ----------
AWS_REGION = get_env("AWS_REGION", "us-west-2")
S3_BUCKET = get_env("S3_BUCKET")

AWS_ACCESS_KEY_ID = get_env("AWS_ACCESS_KEY_ID")
AWS_SECRET_ACCESS_KEY = get_env("AWS_SECRET_ACCESS_KEY")


# ---------- S3 Client ----------
_s3_kwargs = {
    "region_name": AWS_REGION,
}

if AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY:
    _s3_kwargs.update(
        {
            "aws_access_key_id": AWS_ACCESS_KEY_ID,
            "aws_secret_access_key": AWS_SECRET_ACCESS_KEY,
        }
    )

s3 = boto3.client("s3", **_s3_kwargs)

# ---------- Helpers ----------
def md_key_for(d: date) -> str:
    return f"Thailand/{d.strftime('%Y_%m_%d')}.md"


def load_md_from_s3(key: str) -> str:
    obj = s3.get_object(Bucket=S3_BUCKET, Key=key)
    return obj["Body"].read().decode("utf-8")


def list_month_objects(prefix: str) -> set[str]:
    """
    Optional UX improvement:
    - Prefetch list of existing md files in the month to show markers.
    """
    keys = set()
    token = None
    while True:
        kwargs = {"Bucket": S3_BUCKET, "Prefix": prefix, "MaxKeys": 1000}
        if token:
            kwargs["ContinuationToken"] = token
        resp = s3.list_objects_v2(**kwargs)
        for it in resp.get("Contents", []):
            keys.add(it["Key"])
        if resp.get("IsTruncated"):
            token = resp.get("NextContinuationToken")
        else:
            break
    return keys


@dataclass(frozen=True)
class MonthView:
    year: int
    month: int


def month_grid(year: int, month: int):
    cal = calendar.Calendar(firstweekday=0)  # Monday start
    return list(cal.monthdatescalendar(year, month))


# ---------- Page ----------
st.set_page_config(
    page_title="Thailand Daily News",
    page_icon="📰",
    layout="wide",
)

st.title("Thailand Daily News")

now_th = datetime.now(TH_TZ).date()

# Session state
if "selected_date" not in st.session_state:
    st.session_state["selected_date"] = now_th
if "month_view" not in st.session_state:
    st.session_state["month_view"] = MonthView(now_th.year, now_th.month)

selected_date: date = st.session_state["selected_date"]
month_view: MonthView = st.session_state["month_view"]

# CSS: highlight today only (bright background + border)
st.markdown(
    """
<style>
/* Make buttons more compact */
div.stButton > button {
  width: 100%;
  padding: 0.35rem 0.25rem;
  border-radius: 0.75rem;
}

/* Today highlight: we mark via data-testid wrapper class */
.today-btn div.stButton > button {
  border: 2px solid rgba(255, 255, 255, 0.6) !important;
  background: rgba(255, 255, 255, 0.20) !important;
  font-weight: 700 !important;
}

/* Dim out non-current-month buttons (disabled look is already there, but keep subtle) */
.dim-btn div.stButton > button {
  opacity: 0.55;
}
</style>
""",
    unsafe_allow_html=True,
)


# ---------- Layout ----------
col_left, col_right = st.columns([1, 2], gap="large")

with col_left:
    st.subheader("カレンダー")

    # Month navigation
    nav1, nav2, nav3 = st.columns([1, 2, 1])
    with nav1:
        if st.button("", key="prev-month"):
            y, m = month_view.year, month_view.month
            if m == 1:
                month_view = MonthView(y - 1, 12)
            else:
                month_view = MonthView(y, m - 1)
            st.session_state["month_view"] = month_view
            st.rerun()

    with nav2:
        st.markdown(f"### {month_view.year}-{month_view.month:02d}")

    with nav3:
        if st.button("", key="next-month"):
            y, m = month_view.year, month_view.month
            if m == 12:
                month_view = MonthView(y + 1, 1)
            else:
                month_view = MonthView(y, m + 1)
            st.session_state["month_view"] = month_view
            st.rerun()

    # Optional: prefetch existing objects for the month to show indicator
    # Prefix example: Thailand/2026_02_
    month_prefix = f"Thailand/{month_view.year}_{month_view.month:02d}_"
    try:
        existing_keys = list_month_objects(prefix=month_prefix)
    except Exception:
        existing_keys = set()

    # Weekday header
    dow = ["Mon", "Tue", "Wed", "Thu", "Fri", "Sat", "Sun"]
    hdr = st.columns(7)
    for i, d in enumerate(dow):
        hdr[i].markdown(f"**{d}**")

    weeks = month_grid(month_view.year, month_view.month)

    for w in weeks:
        cols = st.columns(7)
        for i, d in enumerate(w):
            is_current_month = (d.month == month_view.month)
            is_today = (d == now_th)

            # Show an indicator if file exists
            key = md_key_for(d)
            has_file = key in existing_keys

            label = str(d.day)
            if has_file and is_current_month:
                label = f"{label}"  # dot marker

            # Use containers with class to control CSS per cell
            wrapper_class = []
            if is_today:
                wrapper_class.append("today-btn")
            if not is_current_month:
                wrapper_class.append("dim-btn")

            # Streamlit doesn't allow per-button class directly, so wrap in HTML
            # that scopes via CSS selectors above.
            with cols[i]:
                st.markdown(f'<div class="{" ".join(wrapper_class)}">', unsafe_allow_html=True)
                clicked = st.button(
                    label,
                    key=f"day-{month_view.year}-{month_view.month}-{d.isoformat()}",
                    disabled=not is_current_month,
                )
                st.markdown("</div>", unsafe_allow_html=True)

                if clicked:
                    st.session_state["selected_date"] = d
                    st.rerun()

    st.divider()
    st.caption("• が付いている日はMarkdownがS3に存在します(推定)。")

with col_right:
    st.subheader(f"記事: {selected_date.strftime('%Y-%m-%d')}")
    key = md_key_for(selected_date)

    try:
        md = load_md_from_s3(key)
        st.markdown(md)
    except ClientError as e:
        code = e.response.get("Error", {}).get("Code", "")
        if code in ("NoSuchKey", "404", "NotFound"):
            st.info(f"この日のファイルがまだありません: `{key}`")
        else:
            st.error(f"S3取得エラー: {e}")
    except Exception as e:
        st.error(f"予期しないエラー: {e}")

実際の挙動

適切に取得できているようです!
記事のリンクも見ましたが、24時間以内に発信された記事を参照できていました。
image.png

最後に

今回は、タイのニュースを自動で収集・要約・保存する仕組みを、AWSとBedrockを使って構築してみました。

まだLINE通知など未実装の部分もありますし、要約精度や対象メディアの拡張など、改善できる余地は多くあります。
今後は、実際の運用を通して調整しながら、より実用的な形に育てていく予定です。

同じように「海外ニュースを効率よく追いたい」「情報収集を自動化したい」と考えている方の参考になれば幸いです。

1
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
1
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?