FlaskとSQLAlchemyで作るPostgreSQLを使ったシンプルなキューの実装
はじめに
アプリケーションで非同期処理を行う際、RabbitMQのようなキューは大変便利です。
ただし、今回プライベートで作成していたアプリケーションでは、シンプルな機能だけで十分でした(可視性タイムアウトやACKなどの高度な機能は必要ありません)。
また、構成を極力簡素に保ちたいため、別途キュー用のミドルウェアを設置することは避けたいと考えました。
そこで、今回はDBにシンプルなキューの役割を担ってもらうことにしました。
調べてみると、Postgres の機能の一つに SKIP LOCKED という機能がありこれを利用することでキューのようなことを実現できそうです。※1
今回は、Flask + SQLAlchemy + PostgreSQL を利用してシンプルなHTTPベースのキューを実装してみました。
利用するもの
- Python
- Flask
- SQLAlchemy
- PostgreSQL -> バージョン 15.5
ざっくり仕組み ※2
PostgreSQLの機能に SKIP LOCKED という機能があります。 これは、ロックできない行をスキップして取得することができるそうです。 SELECT してロックする際に SKIP LOCKED を利用することで、既にロックされた行を無視して取得できます。
実装
- SQLAlchemyでモデルを定義
- メッセージを登録する処理を追加
- メッセージを取得する処理を追加
- Flaskでエンドポイントを定義
- 確認
SQLAlchemyでモデルを定義
まずはモデルを定義します。
messageカラムにJSONで値を入れることができます。
レコードが増える度にidをインクリメントしていきます。
※ 取得する際はidが大きい値から取得します
from sqlalchemy import Column, Integer, JSON class MQ(Base): __tablename__ = "mq" id = Column(Integer, primary_key=True) message = Column(JSON) created_at = Column(DateTime, default=datetime.utcnow)
メッセージを登録する処理追加
メッセージを登録する処理を add_message() に実装します。 先ほど実装した MQというモデルに、messageを登録します
class Query: def __init__( self, DATABASE_URL: str = "postgresql+psycopg2://postgres:postgres@localhost:5432/tegami", ) -> None: engine = create_engine(DATABASE_URL) Session = sessionmaker(bind=engine) self.session = Session() Base.metadata.create_all(engine) def add_message(self, message: dict): model_message = MQ(message=message) self.session.add(model_message) self.session.commit()
メッセージを取得する処理追加
ロックを取得できるものから最大のIDに対してロックをかける処理を行っています。 SQLAlchemyでは with_for_update(skip_locked=True) を利用することで既にロックされているレコードを無視できます。
def get_message(self):
try:
message = (
self.session.query(MQ)
.order_by(MQ.id.asc())
.with_for_update(skip_locked=True)
.limit(1)
.one()
)
self.session.delete(message)
self.session.commit()
return message.message
except Exception as e:
self.session.rollback()
return None
ちなみに、上記コードが実行された際のSQLは下記でした。
SELECT mq.id AS mq_id, mq.message AS mq_message, mq.created_at AS mq_created_at FROM mq ORDER BY mq.id ASC LIMIT %(param_1)s FOR UPDATE SKIP LOCKED
Flaskでエンドポイント定義
/api/message に対してPOSTすることでキューにメッセージを登録することができます。
さらに /api/messageに対してGETすることでメッセージを取得することができます
@app.route("/api/message", methods=["POST"]) def add_message(): req = request.get_json() query.add_message(message=req) return "OK",200 @app.route("/api/message") def get_message(): message = query.get_message() if message: return message,200 else: return "",501
確認
キューにメッセージを積みたい場合
curl -X POST -H "Content-Type: application/json" -d '{"text":"test"}' http://127.0.0.1:8000/api/message
キューからメッセージを取得したい場合
curl http://127.0.0.1:8000/api/message
最後に
Flask + SQLAlchemy + PostgreSQL を使ってシンプルなキューを実装してみました。
MQモデルやモデルへのクエリ処理を修正することでさらに複雑な処理を実装できそうです。