#!/usr/bin/env python3.12 import json import paho.mqtt.client as mqtt import psycopg from psycopg.types.json import Jsonb MQTT_HOST = "192.168.1.20" MQTT_PORT = 1883 MQTT_TOPIC = "sensor/uno01/environment" PG_DSN = ( "host=192.168.1.23" "dbname=mqtt " "user=mqtt " ) def on_connect(client, userdata, flags, reason_code, properties): print(f"MQTT connected: {reason_code}") client.subscribe(MQTT_TOPIC) print(f"subscribe: {MQTT_TOPIC}") def on_message(client, userdata, msg): try: text = msg.payload.decode("utf-8") print(f"recv: topic={msg.topic} payload={text}") # JSONなら解析する try: data = json.loads(text) except json.JSONDecodeError: # 普通の文字列だった場合もJSONBへ保存できる形にする data = { "value": text } # JSONオブジェクトでなければ包む if not isinstance(data, dict): data = { "value": data } device = data.get("device", msg.topic) temperature = data.get("temperature") humidity = data.get("humidity") with psycopg.connect(PG_DSN) as conn: with conn.cursor() as cur: cur.execute( """ INSERT INTO sensor_log (device, temperature, humidity, payload) VALUES (%s, %s, %s, %s) """, ( device, temperature, humidity, Jsonb(data) ) ) print("inserted.") except Exception as e: print(f"ERROR: {e}") client = mqtt.Client( mqtt.CallbackAPIVersion.VERSION2, client_id="mqtt2postgres" ) client.on_connect = on_connect client.on_message = on_message client.connect(MQTT_HOST, MQTT_PORT, 60) client.loop_forever()