You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用Python监听Postgres数据库,实现数据变更实时推送

Hey there! Let's break down how you can listen for Postgres database changes and push that data to your target system—perfect for integrating your ERP with another platform. Here are the most reliable approaches I've used in real-world projects:

1. Use PostgreSQL's Native LISTEN/NOTIFY Mechanism

This is the lightest, built-in way to get change notifications. It works by triggering a NOTIFY event whenever data changes, then having a client listen for those events.

How to set it up:

  • First, create a trigger function that sends a notification when a table is modified. For example, if you're tracking changes to an orders table:
CREATE OR REPLACE FUNCTION notify_order_change()
RETURNS TRIGGER AS $$
BEGIN
  IF TG_OP = 'INSERT' THEN
    PERFORM pg_notify('order_changes', json_build_object(
      'operation', 'INSERT',
      'data', row_to_json(NEW)
    )::text);
  ELSIF TG_OP = 'UPDATE' THEN
    PERFORM pg_notify('order_changes', json_build_object(
      'operation', 'UPDATE',
      'old_data', row_to_json(OLD),
      'new_data', row_to_json(NEW)
    )::text);
  ELSIF TG_OP = 'DELETE' THEN
    PERFORM pg_notify('order_changes', json_build_object(
      'operation', 'DELETE',
      'data', row_to_json(OLD)
    )::text);
  END IF;
  RETURN NULL;
END;
$$ LANGUAGE plpgsql;
  • Attach this trigger to your target table:
CREATE TRIGGER trigger_order_change
AFTER INSERT OR UPDATE OR DELETE ON orders
FOR EACH ROW EXECUTE FUNCTION notify_order_change();
  • Then, write a client (e.g., Python with psycopg2) to listen for these notifications and push data to your target system:
import psycopg2
import psycopg2.extensions
import json

def listen_for_changes():
    conn = psycopg2.connect("dbname=erp_db user=your_user password=your_pass host=localhost")
    conn.set_isolation_level(psycopg2.extensions.ISOLATION_LEVEL_AUTOCOMMIT)
    cur = conn.cursor()
    cur.execute("LISTEN order_changes;")
    
    print("Listening for order changes...")
    while True:
        conn.poll()
        while conn.notifies:
            notify = conn.notifies.pop(0)
            change_data = json.loads(notify.payload)
            # Push this data to your target system here
            print(f"Processing {change_data['operation']} event: {change_data}")
            # Example: call an API or send to a message queue

if __name__ == "__main__":
    listen_for_changes()

Pros & Cons:

  • ✅ Lightweight, no external dependencies
  • ❌ Not persistent: if your listener goes down, you'll miss events
  • ❌ Only sends the data you explicitly include in the trigger—no automatic full change history
2. Logical Replication (With or Without Debezium)

For production-grade reliability, logical replication is the way to go. It uses Postgres's WAL (Write-Ahead Log) to capture full change details, and you can either use Postgres's native subscriptions or a tool like Debezium to stream changes to your target system.

Option A: Native Logical Replication

  • First, make sure your Postgres config has wal_level = logical (edit postgresql.conf and restart the service)
  • Create a publication that includes the tables you want to track:
CREATE PUBLICATION erp_changes FOR TABLE orders, customers;
  • You can set up a Postgres subscription to another database, but if you need to push to a non-Postgres system, you'll need a consumer that reads the replication stream. Tools like pg_recvlogical can help, but it's more low-level.

Debezium is an open-source tool that integrates with Postgres to capture change data streams (CDC) and push them to message brokers like Kafka, which you can then consume to send data to your target system.

  • Debezium connects to your Postgres database, reads the WAL, and emits events with full old/new data, operation type, and metadata.
  • Once you have Debezium set up with Kafka, you can write a simple consumer (e.g., Java, Python) to take these events and push them to your target system's API or database.

Pros & Cons:

  • ✅ Persistent: no data loss if your consumer goes down (Kafka stores events)
  • ✅ Captures full change history, including old and new values
  • ❌ Requires setting up additional components (Kafka, Debezium)
  • ❌ Needs Postgres configuration changes and appropriate permissions
3. Trigger + Intermediate Queue/Table

If you need more flexibility, you can use triggers to write change data to an intermediate log table or a message queue (like RabbitMQ/Redis), then have a separate service read from that queue and push to the target system.

Example with a Log Table:

  • Create a log table to store changes:
CREATE TABLE change_log (
    id SERIAL PRIMARY KEY,
    table_name TEXT NOT NULL,
    operation TEXT NOT NULL,
    old_data JSONB,
    new_data JSONB,
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
  • Modify your trigger function to insert into this table instead of sending a notification:
CREATE OR REPLACE FUNCTION log_change()
RETURNS TRIGGER AS $$
BEGIN
  INSERT INTO change_log (table_name, operation, old_data, new_data)
  VALUES (TG_TABLE_NAME, TG_OP, row_to_json(OLD)::JSONB, row_to_json(NEW)::JSONB);
  RETURN NULL;
END;
$$ LANGUAGE plpgsql;
  • Then, write a service that periodically queries change_log (or uses LISTEN/NOTIFY to get alerts) and processes the entries, pushing them to your target system.

Pros & Cons:

  • ✅ Flexible: you can batch process changes, retry failed pushes, or filter data before sending
  • ✅ Acts as a buffer if the target system is temporarily unavailable
  • ❌ Adds complexity (managing the log table/queue, cleaning up old entries)
Best Practices to Keep in Mind
  • Idempotency: Make sure your target system can handle duplicate events (e.g., use unique IDs to skip already processed data)
  • Filter Unnecessary Changes: Only track tables/columns that your target system needs to reduce load
  • Fault Tolerance: Add retry logic to your push service, and make sure listeners/consumers can reconnect automatically if the connection drops
  • Performance: Avoid heavy logic in triggers—keep them lightweight, and offload processing to your consumer service

内容的提问来源于stack exchange,提问作者Cristopher Werner

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 07:07:59