2026年8月31日 星期一

Python一下:用SQLAlchemy與PyMySQL重構智慧養殖資料存取層

《Python一下:從風土資料到智慧生活》第 28 篇

用SQLAlchemy與PyMySQL重構智慧養殖資料存取層

讓資料表成為Python物件,讓交易有清楚邊界,也讓SQLite與MySQL共用大部分程式。

CH13 關聯式資料庫SQLAlchemy 2.xORM
學習目標
完成本篇後,你能說明Engine、Model、Session與Repository的責任;使用SQLAlchemy 2.x宣告關聯模型;以ORM新增、查詢及更新資料;正確控制交易與連線池;並以SQLite進行快速測試、以PyMySQL連接MySQL。

一、為何要增加SQLAlchemy這一層?

直接使用PyMySQL完全可行,但當SQL散落在API、排程、分析腳本與管理頁面中,欄名修改、交易處理及測試會逐漸困難。SQLAlchemy提供一致的資料存取介面;ORM則把資料表的列映射成Python物件。

元件主要責任水井村例子
Engine管理資料庫方言、連線與連線池MySQL正式環境、SQLite課堂測試
Model描述資料表、欄位與關聯Site、Sensor、Reading
Session追蹤物件及管理一個工作單元新增一次巡查紀錄並提交
Repository集中常用資料存取規則查最新水溫、建立場域
ORM不是不用學SQL:SQLAlchemy最後仍會產生SQL。若不了解JOIN、索引、NULL與交易,就可能寫出「看起來很Python、實際查詢很慢」的程式。

二、安裝SQLAlchemy與PyMySQL

import subprocess
import sys


subprocess.run(
    [
        sys.executable, "-m", "pip", "install",
        "SQLAlchemy>=2.0,<3.0", "PyMySQL>=1.1,<2.0",
    ],
    check=True,
)

實際專案應使用虛擬環境並鎖定經測試的版本。套件升級前先跑測試,不在正式環境臨時更新。

三、建立資料庫網址,但不把密碼寫進程式

import os
from sqlalchemy import URL


def mysql_url_from_environment():
    password = os.getenv("SHUIJING_DB_PASSWORD", "")
    if not password:
        raise RuntimeError("尚未設定SHUIJING_DB_PASSWORD")

    return URL.create(
        drivername="mysql+pymysql",
        username=os.getenv("SHUIJING_DB_USER", "shuijing_app"),
        password=password,
        host=os.getenv("SHUIJING_DB_HOST", "127.0.0.1"),
        port=int(os.getenv("SHUIJING_DB_PORT", "3306")),
        database=os.getenv("SHUIJING_DB_NAME", "shuijing_demo"),
        query={"charset": "utf8mb4"},
    )

URL.create()會處理密碼中的特殊字元。不要print完整URL,因為它可能包含憑證;Colab請使用Secrets,GitHub則使用平台的Secret設定。

四、建立Engine與連線池

from sqlalchemy import create_engine


engine = create_engine(
    mysql_url_from_environment(),
    pool_pre_ping=True,
    pool_recycle=1800,
    pool_size=5,
    max_overflow=5,
    pool_timeout=10,
    echo=False,
)

print(engine.url.render_as_string(hide_password=True))

pool_pre_ping可在借出連線前檢查是否仍有效;連線池大小必須配合MySQL上限、網站程序數與實際流量計算,不可每個程序都任意設很大。

五、先宣告共同Model基底

from sqlalchemy.orm import DeclarativeBase


class Base(DeclarativeBase):
    pass

本篇採SQLAlchemy 2.x的型別化宣告方式。每個Model都繼承Base,後續可由metadata取得所有資料表描述。

六、建立Site、Sensor與Reading模型

from __future__ import annotations

from datetime import datetime
from typing import Optional

from sqlalchemy import (
    BigInteger, Boolean, DateTime, Float, ForeignKey,
    Index, Integer, String, UniqueConstraint,
)
from sqlalchemy.orm import Mapped, mapped_column, relationship


ID_TYPE = BigInteger().with_variant(Integer, "sqlite")


class Site(Base):
    __tablename__ = "sites"

    site_id: Mapped[int] = mapped_column(
        ID_TYPE, primary_key=True, autoincrement=True
    )
    site_code: Mapped[str] = mapped_column(
        String(30), unique=True, nullable=False
    )
    display_name: Mapped[str] = mapped_column(String(100))
    active: Mapped[bool] = mapped_column(Boolean, default=True)
    sensors: Mapped[list[Sensor]] = relationship(
        back_populates="site"
    )


class Sensor(Base):
    __tablename__ = "sensors"

    sensor_id: Mapped[int] = mapped_column(
        ID_TYPE, primary_key=True, autoincrement=True
    )
    sensor_code: Mapped[str] = mapped_column(
        String(50), unique=True, nullable=False
    )
    site_id: Mapped[int] = mapped_column(
        ForeignKey("sites.site_id", ondelete="RESTRICT")
    )
    kind: Mapped[str] = mapped_column(String(30))
    unit: Mapped[str] = mapped_column(String(20))
    site: Mapped[Site] = relationship(back_populates="sensors")
    readings: Mapped[list[Reading]] = relationship(
        back_populates="sensor"
    )


class Reading(Base):
    __tablename__ = "readings"
    __table_args__ = (
        UniqueConstraint(
            "sensor_id", "observed_at",
            name="uq_sensor_time",
        ),
        Index("idx_readings_time", "observed_at"),
    )

    reading_id: Mapped[int] = mapped_column(
        ID_TYPE, primary_key=True, autoincrement=True
    )
    sensor_id: Mapped[int] = mapped_column(
        ForeignKey("sensors.sensor_id", ondelete="RESTRICT")
    )
    observed_at: Mapped[datetime] = mapped_column(
        DateTime(timezone=False)
    )
    value: Mapped[Optional[float]] = mapped_column(Float)
    quality: Mapped[str] = mapped_column(String(10))
    note: Mapped[str] = mapped_column(String(255), default="")
    sensor: Mapped[Sensor] = relationship(
        back_populates="readings"
    )

ID_TYPE讓MySQL使用BIGINT,SQLite測試則使用可正確連結ROWID自動編號的INTEGER。Python型別提示協助編輯器及靜態檢查,但不會自動完成所有資料驗證;quality允許值、時間規則及數值合理性仍應在應用層和資料庫層共同約束。

七、建立資料表:教材可以,正式系統要遷移

Base.metadata.create_all(engine)

create_all()適合初學、測試及全新資料庫,但不會安全地替既有資料表改欄位。正式系統應使用Alembic建立可審查、可追蹤、可回復的結構遷移版本。

八、用Session新增一座場域

from sqlalchemy.orm import Session


with Session(engine) as session:
    with session.begin():
        site = Site(
            site_code="SITE-01",
            display_name="示範池A",
        )
        session.add(site)

print("場域新增完成")

session.begin()區塊正常結束時提交,發生例外時回滾。離開Session後不要任意依賴尚未載入的關聯資料。

九、一次建立場域、感測器與讀值

from datetime import datetime


with Session(engine) as session:
    with session.begin():
        site = Site(
            site_code="SITE-02",
            display_name="示範池B",
        )
        sensor = Sensor(
            sensor_code="TEMP-02",
            kind="temperature",
            unit="°C",
        )
        sensor.readings.append(
            Reading(
                observed_at=datetime(2026, 8, 31, 9, 0),
                value=29.1,
                quality="valid",
            )
        )
        site.sensors.append(sensor)
        session.add(site)

ORM會依關聯順序寫入三張表。教材時間為示例;正式專案應統一以UTC或明確規則寫入,且不可把無時區時間與有時區時間混用。

十、用select()查詢,而不是舊式Query

from sqlalchemy import select


statement = (
    select(Site)
    .where(Site.active.is_(True))
    .order_by(Site.site_code)
)

with Session(engine) as session:
    sites = session.scalars(statement).all()
    for site in sites:
        print(site.site_code, site.display_name)

select(Site)回傳Site物件;如果選擇個別欄位,結果形態會不同。先確認自己需要物件、欄位列,還是單一純量。

十一、JOIN查詢最新有效水溫

statement = (
    select(
        Site.site_code,
        Site.display_name,
        Sensor.sensor_code,
        Reading.observed_at,
        Reading.value,
        Sensor.unit,
    )
    .join(Site.sensors)
    .join(Sensor.readings)
    .where(
        Site.site_code == "SITE-02",
        Sensor.kind == "temperature",
        Reading.quality == "valid",
    )
    .order_by(
        Reading.observed_at.desc(),
        Reading.reading_id.desc(),
    )
    .limit(1)
)

with Session(engine) as session:
    latest = session.execute(statement).mappings().first()
    print(dict(latest) if latest else "尚無有效資料")

SQLAlchemy會把比較值轉成參數,而不是直接拼進SQL。若需要動態欄位或排序,仍必須使用程式內允許清單。

十二、彙總各場域的有效資料

from sqlalchemy import and_, func


statement = (
    select(
        Site.site_code,
        func.count(Reading.reading_id).label("valid_count"),
        func.round(func.avg(Reading.value), 2).label("average"),
    )
    .outerjoin(
        Sensor,
        and_(
            Sensor.site_id == Site.site_id,
            Sensor.kind == "temperature",
        ),
    )
    .outerjoin(
        Reading,
        and_(
            Reading.sensor_id == Sensor.sensor_id,
            Reading.quality == "valid",
        ),
    )
    .group_by(Site.site_id, Site.site_code)
    .order_by(Site.site_code)
)

with Session(engine) as session:
    for row in session.execute(statement).mappings():
        print(dict(row))

把有效資料條件放在LEFT JOIN的ON子句,才能保留沒有有效資料的場域。若放到WHERE,查詢可能悄悄變成只剩有資料的場域。

十三、避免N+1查詢

逐一讀取每座場域的sensors,可能額外送出很多SQL。需要一起使用關聯時,可明確預載:

from sqlalchemy.orm import selectinload


statement = (
    select(Site)
    .options(selectinload(Site.sensors))
    .order_by(Site.site_code)
)

with Session(engine) as session:
    sites = session.scalars(statement).all()
    for site in sites:
        print(site.site_code, len(site.sensors))

selectinload通常以第二個IN查詢載入集合,避免每座場域各查一次。不要為了省查詢而預載所有歷史讀值,資料量可能非常大。

十四、建立輸入驗證函式

from datetime import datetime


def build_reading(observed_at, value, quality, note=""):
    if quality not in {"valid", "unknown", "invalid"}:
        raise ValueError("不支援的資料品質")
    if not isinstance(observed_at, datetime):
        raise TypeError("observed_at必須是datetime")
    if quality == "unknown":
        value = None
    elif value is None:
        raise ValueError("valid或invalid必須保留原始值")
    else:
        value = float(value)

    return Reading(
        observed_at=observed_at,
        value=value,
        quality=quality,
        note=str(note)[:255],
    )
資料品質先於圖表:unknown表示沒有可信數值,invalid表示保留原始值但不宜納入一般統計。不要把兩者都轉成0,否則平均值與告警會被扭曲。

十五、Repository:集中常用存取規則

class SiteRepository:
    def __init__(self, session):
        self.session = session

    def get_by_code(self, site_code):
        statement = select(Site).where(
            Site.site_code == site_code
        )
        return self.session.scalar(statement)

    def add(self, site_code, display_name):
        if self.get_by_code(site_code) is not None:
            raise ValueError("場域代碼已存在")
        site = Site(
            site_code=site_code,
            display_name=display_name,
        )
        self.session.add(site)
        return site


with Session(engine) as session:
    with session.begin():
        repository = SiteRepository(session)
        repository.add("SITE-03", "示範池C")

Repository不應偷偷commit;由外層服務決定交易邊界,才能把多個Repository操作包在同一交易中。

十六、服務層:把場域規則與資料存取分開

def register_sensor(session, site_code, sensor_code, kind, unit):
    repository = SiteRepository(session)
    site = repository.get_by_code(site_code)
    if site is None:
        raise ValueError("場域不存在")

    duplicate = session.scalar(
        select(Sensor).where(Sensor.sensor_code == sensor_code)
    )
    if duplicate is not None:
        raise ValueError("感測器代碼已存在")

    sensor = Sensor(
        sensor_code=sensor_code,
        kind=kind,
        unit=unit,
        site=site,
    )
    session.add(sensor)
    return sensor


with Session(engine) as session:
    with session.begin():
        register_sensor(
            session, "SITE-03", "TEMP-03",
            "temperature", "°C",
        )

十七、用SQLite快速測試相同Model

from sqlalchemy import create_engine
from sqlalchemy.orm import Session


test_engine = create_engine("sqlite+pysqlite:///:memory:")
Base.metadata.create_all(test_engine)

with Session(test_engine) as session:
    with session.begin():
        repository = SiteRepository(session)
        repository.add("TEST-01", "測試場域")

with Session(test_engine) as session:
    saved = session.scalar(
        select(Site).where(Site.site_code == "TEST-01")
    )
    assert saved is not None
    assert saved.display_name == "測試場域"

print("測試通過")

SQLite測試快速,但不能完全代表MySQL:型別、排序規則、鎖定、時區、函式與並行行為可能不同。因此還需要少量MySQL整合測試。

十八、Session常見錯誤

錯誤改善方式
全站共用一個Session每個請求或工作單元建立並關閉Session
Repository自行commit由服務層統一決定提交或回滾
離開Session後才讀延遲關聯在Session內預載或轉為輸出資料
直接回傳ORM物件給外部API選擇必要欄位,避免洩漏內部資料
把echo=True留在正式環境使用受控Log,避免敏感參數外洩

十九、USR資料治理:Model不等於公開資料

Model可能包含場域內部欄位,但學生、養殖戶、公開儀表板與研究資料集需要不同視圖。應在服務層或API輸出層做角色授權與欄位選擇,使用匿名場域代碼,避免輸出姓名、電話、精確座標、設備Token與連線資訊。

def public_site_summary(site, latest_value=None):
    return {
        "site_code": site.site_code,
        "display_name": site.display_name,
        "latest_temperature": latest_value,
    }


# 不輸出內部主鍵、帳號、精確位置或設備憑證

二十、課堂挑戰

挑戰A|基礎:為Reading增加received_at欄位,說明observed_at與received_at的差別。
挑戰B|進階:建立ReadingRepository,完成「批次新增」及「查每支感測器最新有效讀值」,但不要在Repository內commit。
挑戰C|USR場域:設計公開版、養殖戶版及維運版三種輸出資料,逐欄說明誰能看、為何需要,以及保存多久。

二十一、用AI協助審查ORM,而不是盲目產碼

請擔任SQLAlchemy 2.x程式碼審查助教。
系統使用MySQL、PyMySQL,模型為Site、Sensor、Reading。
請檢查:關聯方向、外鍵、唯一鍵、NULL、索引、時區、
Session生命週期、交易邊界、N+1查詢、連線池與敏感資料輸出。
不要改用舊式session.query,也不要要求真實帳密。
請把問題分成「會造成資料錯誤」「效能風險」「維護性」三類,
每項提供最小修正範例與可執行測試。

沒有留言:

張貼留言