0
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?

psycopg(psycopg3)でPostgreSQLに繋いでみた

0
Last updated at Posted at 2026-07-05

概要

  1. PostgreSQLに接続
  2. コネクションプール(AsyncConnectionPool)の作成
  3. SELECT NOW()で現在時刻取得
  4. pg_stat_activityテーブルから現在の状況取得

必要なライブラリ

# tzdataはPostgreSQLから送られてくるタイムゾーン名「Asia/Tokyo」をUTC+09:00と判定するために必要
pip install psycopg psycopg_pool tzdata

ファルダ構成

.
├── log(ログ出力用フォルダ)
│   ├── application.2026-07-04.log
│   └── application.log
├── log_setting.py
├── main.py
└── mymodule
    ├── __init__.py
    ├── dao.py
    └── db.py

ソースコード

main.py

main.py
from typing import Annotated
from csv import QUOTE_NOTNULL, DictWriter
from dataclasses import asdict, fields
from io import StringIO
from pathlib import Path
from psycopg.rows import class_row, dict_row, DictRow
from psycopg.sql import Identifier, SQL, Literal, Placeholder, Composed
from selectors import SelectSelector
from asyncio import run, SelectorEventLoop, sleep
from re import sub
from logging import INFO, getLogger, DEBUG
from mymodule import PgStatActivity, pool
from log_setting import create_log_setting

DataclassAnnotated = Annotated[object, "dataclass"]


async def get_now() -> str | None:

    logger.info('現在時刻の取得')

    try:
        async with pool.connection() as conn:
            await conn.set_read_only(True)

            async with conn.cursor(row_factory=dict_row) as cur:
                await cur.execute('SELECT now() as current', prepare=True)
                row: DictRow | None = await cur.fetchone()  # SQL実行
                logger.info(row)
                if row is not None:
                    # DictRow | NoneというOptional型が返ってくるため、Noneでないことをチェックする
                    if (alias := cur.description[0].name) is not None:
                        # カラム名を取得する(name = 'current')
                        logger.info(row[alias])
                        result = row[alias].strftime('%Y/%m/%d %H:%M:%S')
                        logger.info(result)
                        return result
    except Exception:
        logger.exception('')


async def select(field_class, schema_name: str | None = None, table_name: str | None = None, need_fields: list = ['*'], condition_list: list[SQL, Identifier, Literal, Placeholder, Composed] = []) -> list:

    # テーブル名が無い場合、field_classのアッパーキャメルケースからスネークケースを取得
    if table_name is None:
        snake_case = sub("([A-Z])", lambda x: "_" +
                    x.group(1).lower(), field_class.__name__)
        table_name = snake_case[1:] if snake_case[:1] == '_' else snake_case # 先頭の_を除去

    query: Composed = None
    if len(need_fields) == 0 or need_fields.count('*') > 0:
        # 取得対象フィールドの要素数0、または*が含まれている場合、全項目取得
        query = SQL('SELECT * FROM {table}').format(
            table=Identifier(schema_name, table_name) if schema_name is not None else Identifier(table_name))
    else:
        query = SQL('SELECT {need_field} FROM {table}').format(
            need_field=SQL(',').join(map(Identifier, need_fields)),  # 取得対象
            table=Identifier(
                schema_name, table_name) if schema_name is not None else Identifier(table_name)
        )

    if len(condition_list) > 0:
        # WHEREやORDER BYなどの条件追加
        query += Composed(condition_list)

    logger.debug(query)  # psycopgライブラリ内のクエリ表現

    try:
        async with pool.connection() as conn:
            await conn.set_read_only(True)

            async with conn.cursor(row_factory=class_row(field_class)) as cur:
                logger.info(query.as_string(cur))  # 実行予定の実際のSQL文
                # prepare=TrueでDBのprepareステートメントが強制使用される
                await cur.execute(query, prepare=True)
                result = await cur.fetchall()
                logger.debug(result)  # 実行結果
                logger.debug(cur.description)  # 取得したレコードのカラム情報
                return result
    except Exception:
        logger.exception('')

    return []


async def to_csv(rows: list[DataclassAnnotated], fieldnames: list = [], unlimited=False):

    if len(fieldnames) == 0 and len(rows) > 0:
        # CSVのカラム指定fieldnamesがない場合
        # 無制限を指定しない場合は、隠し要素(prefixとして_を持つ変数名)以外の項目を取得
        col = [f.name for f in fields(
            rows[0]) if unlimited or not f.name.startswith('_')]
        logger.debug(col)
        fieldnames = col

    logger.debug(fieldnames)

    with StringIO() as s:
        writer = DictWriter(s, fieldnames=fieldnames,
                            quoting=QUOTE_NOTNULL, lineterminator='\n')
        writer.writeheader()
        for r in rows:
            original_dict = asdict(r)
            writer.writerow(
                {k: original_dict[k] for k in fieldnames if k in original_dict})
        return s.getvalue()  # CSVデータをstr型で取得


async def main():
    try:
        await pool.open()

        # 現在時刻を取得
        await get_now()

        # 全件取得(条件無し)
        rows = await select(PgStatActivity)
        csv = await to_csv(rows)
        logger.info(f'結果CSV\n{csv}')

        need_fields = ['pid', 'application_name',
                       'query_start', 'state', 'state_change', 'query']

        # 条件(application_name={application_name}で設定した名称で検索し、pidで昇順ソートする)
        logger.info('DBコネクションのレコード状態確認')
        rows = await select(PgStatActivity, need_fields=need_fields,
                            condition_list=[SQL(' WHERE '), Identifier('application_name'), SQL(' = '), Literal('test_application'),
                                            SQL(' ORDER BY '), Identifier('pid')])
        csv = await to_csv(rows, need_fields)
        logger.info(f'結果CSV\n{csv}')
        # await sleep(3)

    except Exception:
        logger.exception('')
    finally:
        await pool.close()


if __name__ == '__main__':
    folder_path = str(Path(__file__).parent) + '/log'
    create_log_setting(folder_path, level=DEBUG)
    logger = getLogger(__name__)

    logger.info('start')

    # psycopgには、PythonデフォルトのProactorEventLoopとは互換性が無いため、loop_factoryを置き換える
    # https://www.psycopg.org/psycopg3/docs/advanced/async.html#asynchronous-operations
    def loop_factory(): return SelectorEventLoop(SelectSelector())
    run(main(), loop_factory=loop_factory)

    logger.info('end')

log_setting.py

log_setting.py
from logging import basicConfig, StreamHandler, INFO, getLogger
from logging.handlers import TimedRotatingFileHandler
from os.path import abspath
from pathlib import Path


class MyTimedRotatingFileHandler(TimedRotatingFileHandler):
    def _open(self):
        # 存在しないフォルダを作成
        Path(self.baseFilename).parent.mkdir(parents=True, exist_ok=True)
        return super()._open()


def create_log_setting(folder, filename='application.log', level=INFO):
    # 毎日ログローテートし、最大30件まで残す
    file_handler = MyTimedRotatingFileHandler(
        # midnightは0時にローテーション(ファイル名はapplication.yyyy-MM-dd.log)
        f'{folder}/{filename}', backupCount=30, encoding='utf-8', when='midnight')
    # 拡張子が末尾に移動するように制御
    file_handler.namer = lambda fn: fn.replace(
        Path(filename).suffix, '') + Path(filename).suffix
    # ログのベース設定を定義
    basicConfig(
        level=level,
        format='%(asctime)s [%(levelname)s] %(name)s - %(message)s',
        datefmt='%Y-%m-%d %H:%M:%S',
        handlers=[
            StreamHandler(),  # コンソール出力
            file_handler      # ファイル出力
        ]
    )
    # ログ取得対象のモジュールを追加
    getLogger('psycopg').setLevel(level)
    getLogger('psycopg.pool').setLevel(level)

    # 出力先の表示
    getLogger(__name__).debug(f'ログ出力先: {abspath(folder)}')

mymodule モジュールのソースコード

init.py

__init__.py
from .db import pool
from .dao import PgStatActivity

# 「from mymodule import pool, PgStatActivity」を記述可能にする
__all__ = [pool, PgStatActivity]

dao.py

dao.py
from dataclasses import dataclass, field
from datetime import datetime


# frozen=True: 値の途中改変不能化、listオブジェクトのappendは利用できる
# kw_only=True: PgStatActivity(datid='xxx', ...)を直接書かないとTypeErrorになる。わざと厳格化して挙動を確認する用
# slots=True: 内部的な辞書(__dict__)にオブジェクトが作られなくなるため、メモリ消費が減る
@dataclass(frozen=True, kw_only=True, slots=True)
class PgStatActivity():
    datid: str = None
    datname: str = None
    pid: str = None
    leader_pid: str = None
    usesysid: str = None
    usename: str = None
    application_name: str = None
    client_addr: str = None
    client_hostname: str = None
    client_port: str = None
    backend_start: datetime = None
    xact_start: datetime = None
    query_start: datetime = None
    state_change: datetime = None
    wait_event_type: str = None
    wait_event: str = None
    state: str = None
    backend_xid: str = None
    backend_xmin: str = None
    query_id: str = None
    query: str = None
    backend_type: str = None

    # 隠し要素
    # default_factory=set/list/dict: 空のオブジェクト型を設定できる
    # init=False: インスタンス化時(__init__)に値初期化をスキップ
    # repr=True: print出力対象外化
    # compare=False: dataclass型同士の比較で比較対象外化
    _hidden: str = field(default=None, init=False, compare=False)
    _str_backend_start: str = field(default=None, init=False, compare=False)
    _str_xact_start: str = field(default=None, init=False, compare=False)
    _str_query_start: str = field(default=None, init=False, compare=False)
    _str_state_change: str = field(default=None, init=False, compare=False)

    # __init__の直後(レコードを当てはめた直後を想定)に動作する
    # レコード取得後のカラム同士の計算等に利用する
    def __post_init__(self):
        # frozen=True なので通常の代入はできない
        # object.__setattr__ を使うと代入できる
        object.__setattr__(self, "_hidden", "隠し項目")
        object.__setattr__(self, "_str_backend_start",
                           '' if self.backend_start is None else self.backend_start.strftime('%Y/%m/%d %H:%M:%S'))
        object.__setattr__(self, "_str_xact_start",
                           '' if self.xact_start is None else self.xact_start.strftime('%Y/%m/%d %H:%M:%S'))
        object.__setattr__(self, "_str_query_start",
                           '' if self.query_start is None else self.query_start.strftime('%Y/%m/%d %H:%M:%S'))
        object.__setattr__(self, "_str_state_change",
                           '' if self.state_change is None else self.state_change.strftime('%Y/%m/%d %H:%M:%S'))

db.py

db.py
from psycopg_pool import AsyncConnectionPool

user = 'postgres'  # pg_stat_activityテーブルを参照するには管理者権限が必要
password = 'postgres'
host = 'localhost'
port = 5432
dbname = 'postgres'
application_name = 'test_application'  # pg_stat_activityテーブル内で識別可能にする
pool = AsyncConnectionPool(
    f'postgresql://{user}:{password}@{host}:{port}/{dbname}?application_name={application_name}',
    # max_idle=10,
    min_size=10,
    max_size=50,
    timeout=10,
    open=False,
    # close_returns=True,
    check=AsyncConnectionPool.check_connection,
)

出力ログイメージ

application.log
application.log
2026-07-05 15:09:40 [DEBUG] log_setting - ログ出力先: <<ログフォルダ>>
2026-07-05 15:09:40 [INFO] __main__ - start
2026-07-05 15:09:40 [DEBUG] asyncio - Using selector: SelectSelector
2026-07-05 15:09:40 [INFO] __main__ - 現在時刻の取得
2026-07-05 15:09:40 [INFO] psycopg.pool - connection requested from 'pool-1'
2026-07-05 15:09:40 [INFO] psycopg.pool - connection given by 'pool-1'
2026-07-05 15:09:40 [INFO] __main__ - {'current': datetime.datetime(2026, 7, 5, 15, 9, 40, 624369, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo'))}
2026-07-05 15:09:40 [INFO] __main__ - 2026-07-05 15:09:40.624369+09:00
2026-07-05 15:09:40 [INFO] __main__ - 2026/07/05 15:09:40
2026-07-05 15:09:40 [INFO] psycopg.pool - returning connection to 'pool-1'
2026-07-05 15:09:40 [DEBUG] __main__ - Composed([SQL('SELECT * FROM '), Identifier('pg_stat_activity')])
2026-07-05 15:09:40 [INFO] psycopg.pool - connection requested from 'pool-1'
2026-07-05 15:09:40 [INFO] psycopg.pool - connection given by 'pool-1'
2026-07-05 15:09:40 [INFO] __main__ - SELECT * FROM "pg_stat_activity"
2026-07-05 15:09:40 [DEBUG] __main__ - [PgStatActivity(datid=5, datname='postgres', pid=10752, leader_pid=None, usesysid=10, usename='postgres', application_name='pgAdmin 4 - DB:postgres', client_addr=IPv6Address('::1'), client_hostname=None, client_port=61308, backend_start=datetime.datetime(2026, 7, 5, 13, 6, 39, 237309, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=None, query_start=datetime.datetime(2026, 7, 5, 13, 7, 20, 733806, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), state_change=datetime.datetime(2026, 7, 5, 13, 7, 20, 769392, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), wait_event_type='Client', wait_event='ClientRead', state='idle', backend_xid=None, backend_xmin=None, query_id=None, query="SELECT obj_type, obj_name,\n    pg_catalog.REPLACE(obj_path, '/'||sn.schema_name||'/', '/'||CASE sn.schema_name\n    WHEN 'pg_catalog' THEN 'PostgreSQL カタログ (pg_catalog)'\n    WHEN 'pgagent' THEN 'pgAgent ジョブスケジューラー (pgagent)'\n    WHEN 'information_schema' THEN 'ANSI (information_schema)'\n    ELSE sn.schema_name\n    END||'/') AS obj_path,\n    schema_name, show_node, other_info,\n    CASE\n        WHEN sn.schema_name IN ('pg_catalog', 'pgagent', 'information_schema') THEN\n            CASE WHEN CASE\n    WHEN sn.schema_name = ANY('{information_schema}')\n        THEN false\n    ELSE true END THEN 'D' ELSE 'O' END\n        ELSE 'N'\n    END AS catalog_level\nFROM (\n    SELECT\n    CASE\n        WHEN c.relkind = 'S' THEN 'sequence'\n        WHEN c.relkind = 'v' THEN 'view'\n        WHEN c.relkind = 'm' THEN 'mview'\n        ELSE 'should not happen'\n    END::text AS obj_type, c.relname AS obj_name,\n    ':schema.'|| n.oid || ':/' || n.nspname || '/' ||\n    CASE\n        WHEN c.relkind = 'S' THEN ':seque", backend_type='client backend', _hidden='隠し項目', _str_backend_start='2026/07/05 13:06:39', _str_xact_start='', _str_query_start='2026/07/05 13:07:20', _str_state_change='2026/07/05 13:07:20'), PgStatActivity(datid=5, datname='postgres', pid=22864, leader_pid=None, usesysid=10, usename='postgres', application_name='test_application', client_addr=IPv6Address('::1'), client_hostname=None, client_port=55808, backend_start=datetime.datetime(2026, 7, 5, 15, 9, 40, 568558, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=datetime.datetime(2026, 7, 5, 15, 9, 40, 627031, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), query_start=datetime.datetime(2026, 7, 5, 15, 9, 40, 631182, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), state_change=datetime.datetime(2026, 7, 5, 15, 9, 40, 631183, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), wait_event_type=None, wait_event=None, state='active', backend_xid=None, backend_xmin='775', query_id=None, query='SELECT * FROM "pg_stat_activity"', backend_type='client backend', _hidden='隠し項目', _str_backend_start='2026/07/05 15:09:40', _str_xact_start='2026/07/05 15:09:40', _str_query_start='2026/07/05 15:09:40', _str_state_change='2026/07/05 15:09:40'), PgStatActivity(datid=16390, datname='test', pid=13588, leader_pid=None, usesysid=10, usename='postgres', application_name='pgAdmin 4 - DB:test', client_addr=IPv6Address('::1'), client_hostname=None, client_port=61324, backend_start=datetime.datetime(2026, 7, 5, 13, 6, 39, 609022, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=None, query_start=datetime.datetime(2026, 7, 5, 13, 6, 58, 402509, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), state_change=datetime.datetime(2026, 7, 5, 13, 6, 58, 402862, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), wait_event_type='Client', wait_event='ClientRead', state='idle', backend_xid=None, backend_xmin=None, query_id=None, query="SELECT\n    nsp.nspname as schema_name,\n    (nsp.nspname = 'pg_catalog' AND EXISTS\n        (SELECT 1 FROM pg_catalog.pg_class WHERE relname = 'pg_class' AND\n            relnamespace = nsp.oid LIMIT 1)) OR\n    (nsp.nspname = 'pgagent' AND EXISTS\n        (SELECT 1 FROM pg_catalog.pg_class WHERE relname = 'pga_job' AND\n            relnamespace = nsp.oid LIMIT 1)) OR\n    (nsp.nspname = 'information_schema' AND EXISTS\n        (SELECT 1 FROM pg_catalog.pg_class WHERE relname = 'tables' AND\n            relnamespace = nsp.oid LIMIT 1)) AS is_catalog,\n    CASE\n    WHEN nsp.nspname = ANY('{information_schema}')\n        THEN false\n    ELSE true END AS db_support\nFROM\n    pg_catalog.pg_namespace nsp\nWHERE\n    nsp.oid = 16391::OID;", backend_type='client backend', _hidden='隠し項目', _str_backend_start='2026/07/05 13:06:39', _str_xact_start='', _str_query_start='2026/07/05 13:06:58', _str_state_change='2026/07/05 13:06:58'), PgStatActivity(datid=5, datname='postgres', pid=17308, leader_pid=None, usesysid=10, usename='postgres', application_name='test_application', client_addr=IPv6Address('::1'), client_hostname=None, client_port=55806, backend_start=datetime.datetime(2026, 7, 5, 15, 9, 40, 563872, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=None, query_start=datetime.datetime(2026, 7, 5, 15, 9, 40, 623242, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), state_change=datetime.datetime(2026, 7, 5, 15, 9, 40, 623291, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), wait_event_type='Client', wait_event='ClientRead', state='idle', backend_xid=None, backend_xmin=None, query_id=None, query='COMMIT', backend_type='client backend', _hidden='隠し項目', _str_backend_start='2026/07/05 15:09:40', _str_xact_start='', _str_query_start='2026/07/05 15:09:40', _str_state_change='2026/07/05 15:09:40'), PgStatActivity(datid=5, datname='postgres', pid=9184, leader_pid=None, usesysid=10, usename='postgres', application_name='test_application', client_addr=IPv6Address('::1'), client_hostname=None, client_port=55807, backend_start=datetime.datetime(2026, 7, 5, 15, 9, 40, 566906, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=None, query_start=datetime.datetime(2026, 7, 5, 15, 9, 40, 625830, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), state_change=datetime.datetime(2026, 7, 5, 15, 9, 40, 625862, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), wait_event_type='Client', wait_event='ClientRead', state='idle', backend_xid=None, backend_xmin=None, query_id=None, query='COMMIT', backend_type='client backend', _hidden='隠し項目', _str_backend_start='2026/07/05 15:09:40', _str_xact_start='', _str_query_start='2026/07/05 15:09:40', _str_state_change='2026/07/05 15:09:40'), PgStatActivity(datid=None, datname=None, pid=21564, leader_pid=None, usesysid=None, usename=None, application_name='', client_addr=None, client_hostname=None, client_port=None, backend_start=datetime.datetime(2026, 7, 4, 20, 38, 0, 408228, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=None, query_start=None, state_change=None, wait_event_type='Activity', wait_event='AutovacuumMain', state=None, backend_xid=None, backend_xmin=None, query_id=None, query='', backend_type='autovacuum launcher', _hidden='隠し項目', _str_backend_start='2026/07/04 20:38:00', _str_xact_start='', _str_query_start='', _str_state_change=''), PgStatActivity(datid=None, datname=None, pid=4852, leader_pid=None, usesysid=10, usename='postgres', application_name='', client_addr=None, client_hostname=None, client_port=None, backend_start=datetime.datetime(2026, 7, 4, 20, 38, 0, 409166, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=None, query_start=None, state_change=None, wait_event_type='Activity', wait_event='LogicalLauncherMain', state=None, backend_xid=None, backend_xmin=None, query_id=None, query='', backend_type='logical replication launcher', _hidden='隠し項目', _str_backend_start='2026/07/04 20:38:00', _str_xact_start='', _str_query_start='', _str_state_change=''), PgStatActivity(datid=None, datname=None, pid=11720, leader_pid=None, usesysid=None, usename=None, application_name='', client_addr=None, client_hostname=None, client_port=None, backend_start=datetime.datetime(2026, 7, 4, 20, 38, 0, 32475, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=None, query_start=None, state_change=None, wait_event_type='Activity', wait_event='IoWorkerMain', state=None, backend_xid=None, backend_xmin=None, query_id=None, query='', backend_type='io worker', _hidden='隠し項目', _str_backend_start='2026/07/04 20:38:00', _str_xact_start='', _str_query_start='', _str_state_change=''), PgStatActivity(datid=None, datname=None, pid=26112, leader_pid=None, usesysid=None, usename=None, application_name='', client_addr=None, client_hostname=None, client_port=None, backend_start=datetime.datetime(2026, 7, 4, 20, 38, 0, 37227, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=None, query_start=None, state_change=None, wait_event_type='Activity', wait_event='IoWorkerMain', state=None, backend_xid=None, backend_xmin=None, query_id=None, query='', backend_type='io worker', _hidden='隠し項目', _str_backend_start='2026/07/04 20:38:00', _str_xact_start='', _str_query_start='', _str_state_change=''), PgStatActivity(datid=None, datname=None, pid=25416, leader_pid=None, usesysid=None, usename=None, application_name='', client_addr=None, client_hostname=None, client_port=None, backend_start=datetime.datetime(2026, 7, 4, 20, 38, 0, 40870, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=None, query_start=None, state_change=None, wait_event_type='Activity', wait_event='IoWorkerMain', state=None, backend_xid=None, backend_xmin=None, query_id=None, query='', backend_type='io worker', _hidden='隠し項目', _str_backend_start='2026/07/04 20:38:00', _str_xact_start='', _str_query_start='', _str_state_change=''), PgStatActivity(datid=None, datname=None, pid=5520, leader_pid=None, usesysid=None, usename=None, application_name='', client_addr=None, client_hostname=None, client_port=None, backend_start=datetime.datetime(2026, 7, 4, 20, 38, 0, 44524, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=None, query_start=None, state_change=None, wait_event_type='Activity', wait_event='CheckpointerMain', state=None, backend_xid=None, backend_xmin=None, query_id=None, query='', backend_type='checkpointer', _hidden='隠し項目', _str_backend_start='2026/07/04 20:38:00', _str_xact_start='', _str_query_start='', _str_state_change=''), PgStatActivity(datid=None, datname=None, pid=27704, leader_pid=None, usesysid=None, usename=None, application_name='', client_addr=None, client_hostname=None, client_port=None, backend_start=datetime.datetime(2026, 7, 4, 20, 38, 0, 45560, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=None, query_start=None, state_change=None, wait_event_type='Activity', wait_event='BgwriterHibernate', state=None, backend_xid=None, backend_xmin=None, query_id=None, query='', backend_type='background writer', _hidden='隠し項目', _str_backend_start='2026/07/04 20:38:00', _str_xact_start='', _str_query_start='', _str_state_change=''), PgStatActivity(datid=None, datname=None, pid=15192, leader_pid=None, usesysid=None, usename=None, application_name='', client_addr=None, client_hostname=None, client_port=None, backend_start=datetime.datetime(2026, 7, 4, 20, 38, 0, 403638, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), xact_start=None, query_start=None, state_change=None, wait_event_type='Activity', wait_event='WalWriterMain', state=None, backend_xid=None, backend_xmin=None, query_id=None, query='', backend_type='walwriter', _hidden='隠し項目', _str_backend_start='2026/07/04 20:38:00', _str_xact_start='', _str_query_start='', _str_state_change='')]
2026-07-05 15:09:40 [DEBUG] __main__ - [<Column 'datid', type: oid (oid: 26)>, <Column 'datname', type: name (oid: 19)>, <Column 'pid', type: int4 (oid: 23)>, <Column 'leader_pid', type: int4 (oid: 23)>, <Column 'usesysid', type: oid (oid: 26)>, <Column 'usename', type: name (oid: 19)>, <Column 'application_name', type: text (oid: 25)>, <Column 'client_addr', type: inet (oid: 869)>, <Column 'client_hostname', type: text (oid: 25)>, <Column 'client_port', type: int4 (oid: 23)>, <Column 'backend_start', type: timestamptz (oid: 1184)>, <Column 'xact_start', type: timestamptz (oid: 1184)>, <Column 'query_start', type: timestamptz (oid: 1184)>, <Column 'state_change', type: timestamptz (oid: 1184)>, <Column 'wait_event_type', type: text (oid: 25)>, <Column 'wait_event', type: text (oid: 25)>, <Column 'state', type: text (oid: 25)>, <Column 'backend_xid', type: xid (oid: 28)>, <Column 'backend_xmin', type: xid (oid: 28)>, <Column 'query_id', type: int8 (oid: 20)>, <Column 'query', type: text (oid: 25)>, <Column 'backend_type', type: text (oid: 25)>]
2026-07-05 15:09:40 [INFO] psycopg.pool - returning connection to 'pool-1'
2026-07-05 15:09:40 [DEBUG] __main__ - ['datid', 'datname', 'pid', 'leader_pid', 'usesysid', 'usename', 'application_name', 'client_addr', 'client_hostname', 'client_port', 'backend_start', 'xact_start', 'query_start', 'state_change', 'wait_event_type', 'wait_event', 'state', 'backend_xid', 'backend_xmin', 'query_id', 'query', 'backend_type']
2026-07-05 15:09:40 [INFO] __main__ - 結果CSV
"datid","datname","pid","leader_pid","usesysid","usename","application_name","client_addr","client_hostname","client_port","backend_start","xact_start","query_start","state_change","wait_event_type","wait_event","state","backend_xid","backend_xmin","query_id","query","backend_type"
"5","postgres","31516",,"10","postgres","test_application","::1",,"49225","2026-07-15 21:18:49.848021+09:00",,,"2026-07-15 21:18:49.892009+09:00","Client","ClientRead","idle",,,,"","client backend"
"5","postgres","28280",,"10","postgres","test_application","::1",,"49224","2026-07-15 21:18:49.846497+09:00",,,"2026-07-15 21:18:49.883554+09:00","Client","ClientRead","idle",,,,"","client backend"
"5","postgres","28952",,"10","postgres","test_application","::1",,"49223","2026-07-15 21:18:49.844381+09:00","2026-07-15 21:18:49.893175+09:00","2026-07-15 21:18:49.899775+09:00","2026-07-15 21:18:49.899776+09:00",,,"active",,"2014",,"SELECT * FROM ""pg_stat_activity""","client backend"
"5","postgres","30472",,"10","postgres","pgAdmin 4 - CONN:6481458","::1",,"62968","2026-07-15 20:37:10.568697+09:00",,,"2026-07-15 20:37:10.598977+09:00","Client","ClientRead","idle",,,,"","client backend"
,,"8480",,,,"",,,,"2026-07-12 16:47:48.963063+09:00",,,,"Activity","AutovacuumMain",,,,,"","autovacuum launcher"
,,"8488",,"10","postgres","",,,,"2026-07-12 16:47:48.962308+09:00",,,,"Activity","LogicalLauncherMain",,,,,"","logical replication launcher"
,,"7984",,,,"",,,,"2026-07-12 16:47:48.791977+09:00",,,,"Activity","IoWorkerMain",,,,,"","io worker"
,,"8028",,,,"",,,,"2026-07-12 16:47:48.799123+09:00",,,,"Activity","IoWorkerMain",,,,,"","io worker"
,,"8052",,,,"",,,,"2026-07-12 16:47:48.804223+09:00",,,,"Activity","IoWorkerMain",,,,,"","io worker"
,,"8084",,,,"",,,,"2026-07-12 16:47:48.809754+09:00",,,,"Activity","CheckpointerMain",,,,,"","checkpointer"
,,"8472",,,,"",,,,"2026-07-12 16:47:48.957926+09:00",,,,"Activity","WalWriterMain",,,,,"","walwriter"
,,"8132",,,,"",,,,"2026-07-12 16:47:48.820262+09:00",,,,"Activity","BgwriterHibernate",,,,,"","background writer"

2026-07-05 15:09:40 [INFO] __main__ - DBコネクションの状態確認
2026-07-05 15:09:40 [DEBUG] __main__ - Composed([SQL('SELECT '), Composed([Identifier('pid'), SQL(','), Identifier('application_name'), SQL(','), Identifier('query_start'), SQL(','), Identifier('state'), SQL(','), Identifier('state_change'), SQL(','), Identifier('query')]), SQL(' FROM '), Identifier('pg_stat_activity'), SQL(' WHERE '), Identifier('application_name'), SQL(' = '), Literal('test_application'), SQL(' ORDER BY '), Identifier('pid')])
2026-07-05 15:09:40 [INFO] psycopg.pool - connection requested from 'pool-1'
2026-07-05 15:09:40 [INFO] psycopg.pool - connection given by 'pool-1'
2026-07-05 15:09:40 [INFO] __main__ - SELECT "pid","application_name","query_start","state","state_change","query" FROM "pg_stat_activity" WHERE "application_name" = 'test_application' ORDER BY "pid"
2026-07-05 15:09:40 [DEBUG] __main__ - [PgStatActivity(datid=None, datname=None, pid=9184, leader_pid=None, usesysid=None, usename=None, application_name='test_application', client_addr=None, client_hostname=None, client_port=None, backend_start=None, xact_start=None, query_start=datetime.datetime(2026, 7, 5, 15, 9, 40, 625830, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), state_change=datetime.datetime(2026, 7, 5, 15, 9, 40, 625862, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), wait_event_type=None, wait_event=None, state='idle', backend_xid=None, backend_xmin=None, query_id=None, query='COMMIT', backend_type=None, _hidden='隠し項目', _str_backend_start='', _str_xact_start='', _str_query_start='2026/07/05 15:09:40', _str_state_change='2026/07/05 15:09:40'), PgStatActivity(datid=None, datname=None, pid=17308, leader_pid=None, usesysid=None, usename=None, application_name='test_application', client_addr=None, client_hostname=None, client_port=None, backend_start=None, xact_start=None, query_start=datetime.datetime(2026, 7, 5, 15, 9, 40, 639682, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), state_change=datetime.datetime(2026, 7, 5, 15, 9, 40, 639682, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), wait_event_type=None, wait_event=None, state='active', backend_xid=None, backend_xmin=None, query_id=None, query='SELECT "pid","application_name","query_start","state","state_change","query" FROM "pg_stat_activity" WHERE "application_name" = \'test_application\' ORDER BY "pid"', backend_type=None, _hidden='隠し項目', _str_backend_start='', _str_xact_start='', _str_query_start='2026/07/05 15:09:40', _str_state_change='2026/07/05 15:09:40'), PgStatActivity(datid=None, datname=None, pid=22864, leader_pid=None, usesysid=None, usename=None, application_name='test_application', client_addr=None, client_hostname=None, client_port=None, backend_start=None, xact_start=None, query_start=datetime.datetime(2026, 7, 5, 15, 9, 40, 633533, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), state_change=datetime.datetime(2026, 7, 5, 15, 9, 40, 633595, tzinfo=zoneinfo.ZoneInfo(key='Asia/Tokyo')), wait_event_type=None, wait_event=None, state='idle', backend_xid=None, backend_xmin=None, query_id=None, query='COMMIT', backend_type=None, _hidden='隠し項目', _str_backend_start='', _str_xact_start='', _str_query_start='2026/07/05 15:09:40', _str_state_change='2026/07/05 15:09:40')]
2026-07-05 15:09:40 [DEBUG] __main__ - [<Column 'pid', type: int4 (oid: 23)>, <Column 'application_name', type: text (oid: 25)>, <Column 'query_start', type: timestamptz (oid: 1184)>, <Column 'state', type: text (oid: 25)>, <Column 'state_change', type: timestamptz (oid: 1184)>, <Column 'query', type: text (oid: 25)>]
2026-07-05 15:09:40 [INFO] psycopg.pool - returning connection to 'pool-1'
2026-07-05 15:09:40 [INFO] __main__ - 結果CSV
"pid","application_name","query_start","state","state_change","query"
9184,"test_application","2026-07-05 15:09:40.625830+09:00","idle","2026-07-05 15:09:40.625862+09:00","COMMIT"
17308,"test_application","2026-07-05 15:09:40.639682+09:00","active","2026-07-05 15:09:40.639682+09:00","SELECT ""pid"",""application_name"",""query_start"",""state"",""state_change"",""query"" FROM ""pg_stat_activity"" WHERE ""application_name"" = 'test_application' ORDER BY ""pid"""
22864,"test_application","2026-07-05 15:09:40.633533+09:00","idle","2026-07-05 15:09:40.633595+09:00","COMMIT"

2026-07-05 15:09:40 [INFO] psycopg.pool - adding new connection to the pool
2026-07-05 15:09:40 [INFO] psycopg.pool - adding new connection to the pool
2026-07-05 15:09:40 [INFO] psycopg.pool - adding new connection to the pool
2026-07-05 15:09:40 [INFO] __main__ - end

つまずきポイント(デフォルト設定で動かないことが多い...)

asyncioライブラリのデフォルトループ

NG

async def main(): ...AsyncConnectionPool(...)
asyncio.run(main()) # デフォルトループを使用

OK
ProactorEventLoop ループをSelectorEventLoopループに差し替える必要がある

async def main(): ...AsyncConnectionPool(...)
asyncio.run(main(), loop_factory=lambda: asyncio.SelectorEventLoop(selectors.SelectSelector())) # SelectorEventLoop使用

コネクションプールのオープンタイミング(アプリ起動時と同時はデフォルト値にも関わらず、非推奨化されている)

NG

# グローバルな宣言でopen=True(デフォルト)を使うのは不可能
pool = psycopg_pool.AsyncConnectionPool(
    conninfo='user=postgres password=postgres',
    open=True # open=TrueのままAsyncConnectionPoolをwith句なしで呼び出すのはエラーで不可能
)

OK
open=Falseが推奨されている

In a future version, the default value for the open parameter might be changed to False.
将来のバージョンでは、open パラメータのデフォルト値が False に変更される可能性があります。

pool = psycopg_pool.AsyncConnectionPool(
    conninfo='user=postgres password=postgres',
    open=False # OK
)

with句を使う方法でopen=Trueは使えるが、起動時に非推奨機能であることを伝えるメッセージが表示される

RuntimeWarning: opening the async pool AsyncConnectionPool in the constructor is deprecated and will not be supported anymore in a future release. Please use await pool.open(), or use the pool as context manager using: async with AsyncConnectionPool(...) as pool: ...

import asyncio
import selectors
from psycopg_pool import AsyncConnectionPool
async def main():
    async with AsyncConnectionPool(conninfo='user=postgres password=postgres', open=True) as pool:
        async with pool.connection() as conn:
            async with conn.cursor() as cur:
                await cur.execute('SELECT VERSION()')
                result = await cur.fetchone()
                print(*result) # 「PostgreSQL 18.4 on x86_64-windows, compiled by msvc-19.44.35227, 64-bit」等

asyncio.run(main(), loop_factory=lambda: asyncio.SelectorEventLoop(selectors.SelectSelector()))

結論はグローバル変数でAsyncConnectionPoolをインスタンス化して、async def内でawait pool.open()/close()を使い、コネクションプールを呼び出すのが正しい

pool = psycopg_pool.AsyncConnectionPool(
    conninfo='user=postgres password=postgres',
    open=False
)

async def main():
    await pool.open() # 実質的に正規手法
    # ...
    await pool.close()

コネクションプールのオープン・クローズの繰り返し

NG

await pool.open()
await pool.close() # コネクションプールを破棄するコマンドです
await pool.open()  # いかなる場合でも再オープン出来ません

OK
一回のオープンで全てのSQLを実行しましょう。

from asyncio import run, SelectorEventLoop
from psycopg_pool import AsyncConnectionPool
from selectors import SelectSelector
from traceback import print_exc


pool = AsyncConnectionPool(
    conninfo='user=postgres password=postgres',
    open=False,
)


async def exec_query(p: AsyncConnectionPool, sql: str):
    async with p.connection() as conn:
        await conn.set_read_only(True)

        print(p.get_stats())  # コネクション状況確認
        async with conn.cursor() as cur:
            await cur.execute(sql)
            result = await cur.fetchone()
            print(*result)
        print(p.get_stats())  # コネクション状況確認


async def main():
    try:
        await pool.open()
        await exec_query(pool, 'SELECT VERSION()')
        await exec_query(pool, 'SELECT current_timestamp')
        await exec_query(pool, 'SELECT 1')
        await exec_query(pool, 'SELECT NOW()')
    except Exception:
        print_exc()
    finally:
        await pool.close()


if __name__ == '__main__':
    run(main(), loop_factory=lambda: SelectorEventLoop(SelectSelector()))

select.pyというファイル名は厳禁

selectorsライブラリに既に存在するため

COPY文を使えば、SQLでもCSVを作成可能

結果の取得方法が独特で、memoryview型で得られるバイト列をデコードする必要がある。
また、COPY文の仕様でヘッダー行にダブルクォーテーションを付けることができない。

main.py 改修後
async def copy_to_csv(field_class, schema_name: str | None = None, table_name: str | None = None, need_fields: list = ['*'], condition_list: list[SQL, Identifier, Literal, Placeholder, Composed] = []) -> str:

    # テーブル名が無い場合、field_classのアッパーキャメルケースからスネークケースを取得
    if table_name is None:
        snake_case = sub("([A-Z])", lambda x: "_" +
                         x.group(1).lower(), field_class.__name__)
        # 先頭の_を除去
        table_name = snake_case[1:] if snake_case[:1] == '_' else snake_case

    query: Composed = None
    if len(need_fields) == 0 or need_fields.count('*') > 0:
        # 取得対象フィールドの要素数0、または*が含まれている場合、全項目取得
        query = SQL('COPY (SELECT * FROM {table})').format(
            table=Identifier(schema_name, table_name) if schema_name is not None else Identifier(table_name))
    else:
        query = SQL('COPY (SELECT {need_field} FROM {table})').format(
            need_field=SQL(',').join(map(Identifier, need_fields)),  # 取得対象
            table=Identifier(
                schema_name, table_name) if schema_name is not None else Identifier(table_name)
        )

    # CSV形式のための設定
    query += SQL(" TO STDOUT WITH (FORMAT CSV, HEADER, NULL '', DELIMITER ',', QUOTE '\"', FORCE_QUOTE *,  ENCODING 'UTF8')")

    logger.debug(query)  # psycopgライブラリ内のクエリ表現

    try:
        async with pool.connection() as conn:
            await conn.set_read_only(True)

            async with conn.cursor() as cur:
                logger.info(query.as_string(cur))  # 実行予定の実際のSQL文

                async with cur.copy(query) as copy:
                    result = ''
                    while row := await copy.read():
                        # memoryview型で取得されるため、デコードする必要がある
                        result += bytes(row).decode("utf-8").strip(' ')
                    logger.debug(result)
                    logger.debug(cur.description)  # 取得したレコードのカラム情報
                    return result

    except Exception:
        logger.exception('')

    return ''


async def main():
    try:
        await pool.open()
        #...中略
        csv = await copy_to_csv(PgStatActivity)
        logger.info(f'COPY TO CSV\n{csv}')
        #...中略
    except Exception:
        logger.exception('')
    finally:
        await pool.close()

0
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
0
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?