用SQLAlchemy與PyMySQL重構智慧養殖資料存取層
讓資料表成為Python物件,讓交易有清楚邊界,也讓SQLite與MySQL共用大部分程式。
完成本篇後,你能說明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 | 集中常用資料存取規則 | 查最新水溫、建立場域 |
二、安裝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],
)十五、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,
}
# 不輸出內部主鍵、帳號、精確位置或設備憑證二十、課堂挑戰
二十一、用AI協助審查ORM,而不是盲目產碼
請擔任SQLAlchemy 2.x程式碼審查助教。
系統使用MySQL、PyMySQL,模型為Site、Sensor、Reading。
請檢查:關聯方向、外鍵、唯一鍵、NULL、索引、時區、
Session生命週期、交易邊界、N+1查詢、連線池與敏感資料輸出。
不要改用舊式session.query,也不要要求真實帳密。
請把問題分成「會造成資料錯誤」「效能風險」「維護性」三類,
每項提供最小修正範例與可執行測試。
沒有留言:
張貼留言