FlaskとSQLAlchemyで作るPostgreSQLを使ったシンプルなキューの実装

はじめに

アプリケーションで非同期処理を行う際、RabbitMQのようなキューは大変便利です。 ただし、今回プライベートで作成していたアプリケーションでは、シンプルな機能だけで十分でした(可視性タイムアウトやACKなどの高度な機能は必要ありません)。
また、構成を極力簡素に保ちたいため、別途キュー用のミドルウェアを設置することは避けたいと考えました。

そこで、今回はDBにシンプルなキューの役割を担ってもらうことにしました。 調べてみると、Postgres の機能の一つに SKIP LOCKED という機能がありこれを利用することでキューのようなことを実現できそうです。※1
今回は、Flask + SQLAlchemy + PostgreSQL を利用してシンプルなHTTPベースのキューを実装してみました。

利用するもの

ざっくり仕組み ※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モデルやモデルへのクエリ処理を修正することでさらに複雑な処理を実装できそうです。

参考