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

如何使用Python捕获Oracle 11g Release 2的数据变更(Insert DML)

How to Capture Insert DML Changes in Oracle 11g Release 2 Using Python

Got it, let's tackle how to capture insert DML changes in Oracle 11g R2 with Python—since you mentioned Java has tools like Streamsets but Python examples are hard to find, I'll break down practical, actionable approaches you can use right away.

Approach 1: Trigger-Based CDC (Simple, Low Overhead for Small Workloads)

This is the most straightforward method for smaller datasets. We'll create a database trigger that writes every new inserted record to a dedicated change log table, then use Python to poll this table for new entries.

Step 1: Set Up the Change Log Table and Trigger

First, run these SQL commands in your Oracle database (you'll need appropriate permissions to create tables and triggers):

-- Create a change log table matching your target table's structure (adjust columns as needed)
CREATE TABLE user_data_changelog (
    changelog_id NUMBER GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    original_id NUMBER, -- Match the primary key of your target table
    username VARCHAR2(50),
    email VARCHAR2(100),
    change_timestamp TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
    operation_type VARCHAR2(10) DEFAULT 'INSERT'
);

-- Create an AFTER INSERT trigger on your target table (replace 'user_data' with your table name)
CREATE OR REPLACE TRIGGER trg_user_data_insert
AFTER INSERT ON user_data
FOR EACH ROW
BEGIN
    INSERT INTO user_data_changelog (original_id, username, email)
    VALUES (:NEW.id, :NEW.username, :NEW.email);
END;
/

Step 2: Python Code to Poll for New Inserts

Use cx_Oracle (Oracle's official Python driver) to connect to the database and fetch new changes. We'll track the last processed changelog_id to avoid reprocessing data:

import cx_Oracle
import time

def get_new_inserts(last_processed_id=0):
    # Replace with your Oracle connection details
    dsn_tns = cx_Oracle.makedsn("your_host", "your_port", service_name="your_service")
    connection = cx_Oracle.connect(user="your_user", password="your_password", dsn=dsn_tns)
    cursor = connection.cursor()

    # Fetch all new insert records since last processed ID
    cursor.execute("""
        SELECT changelog_id, original_id, username, email, change_timestamp
        FROM user_data_changelog
        WHERE changelog_id > :last_id
        ORDER BY changelog_id ASC
    """, last_id=last_processed_id)

    changes = cursor.fetchall()
    cursor.close()
    connection.close()

    return changes

# Example: Poll every 10 seconds for new inserts
last_id = 0
while True:
    new_changes = get_new_inserts(last_id)
    if new_changes:
        print("New insert detected:")
        for change in new_changes:
            print(f"ID: {change[1]}, Username: {change[2]}, Email: {change[3]}, Timestamp: {change[4]}")
            last_id = change[0]  # Update last processed ID
    time.sleep(10)

Approach 2: Oracle LogMiner-Based CDC (Efficient, No Trigger Overhead)

Oracle 11g R2 includes LogMiner, a built-in tool that parses redo and archive logs to extract DML operations. This method avoids adding triggers to your source tables, making it better for larger workloads.

Step 1: Grant Necessary Permissions

Your Oracle user needs these permissions to use LogMiner:

GRANT SELECT_CATALOG_ROLE TO your_user;
GRANT EXECUTE ON DBMS_LOGMNR TO your_user;
GRANT EXECUTE ON DBMS_LOGMNR_D TO your_user;

Step 2: Python Code to Extract Insert Changes via LogMiner

This code starts LogMiner, targets the relevant logs, and queries the V$LOGMNR_CONTENTS view to filter insert operations:

import cx_Oracle

def extract_insert_changes():
    dsn_tns = cx_Oracle.makedsn("your_host", "your_port", service_name="your_service")
    connection = cx_Oracle.connect(user="your_user", password="your_password", dsn=dsn_tns)
    cursor = connection.cursor()

    # Start LogMiner with current redo logs (adjust options as needed)
    cursor.execute("BEGIN DBMS_LOGMNR.START_LOGMNR(OPTIONS => DBMS_LOGMNR.DICT_FROM_ONLINE_CATALOG + DBMS_LOGMNR.CONTINUOUS_MINE); END;")

    # Query for INSERT operations on your target table
    cursor.execute("""
        SELECT SQL_REDO, TIMESTAMP, SEG_OWNER, SEG_NAME
        FROM V$LOGMNR_CONTENTS
        WHERE OPERATION = 'INSERT'
        AND SEG_OWNER = 'YOUR_SCHEMA'  -- Replace with your schema name
        AND SEG_NAME = 'USER_DATA'     -- Replace with your table name
    """)

    insert_changes = cursor.fetchall()
    print("Extracted INSERT operations:")
    for change in insert_changes:
        print(f"Timestamp: {change[1]}, Table: {change[2]}.{change[3]}")
        print(f"SQL: {change[0]}\n")

    # Stop LogMiner
    cursor.execute("BEGIN DBMS_LOGMNR.END_LOGMNR(); END;")

    cursor.close()
    connection.close()

extract_insert_changes()

Notes on LogMiner:

  • You can add specific archive logs instead of using continuous mining by calling DBMS_LOGMNR.ADD_LOGFILE() before starting LogMiner.
  • For long-running CDC, consider running LogMiner in a loop and tracking the last processed SCN (System Change Number) to avoid reprocessing logs.

Approach 3: Oracle GoldenGate (Enterprise-Grade CDC)

If you have access to Oracle GoldenGate (a separate CDC tool), you can use Python to interact with it:

  • Option 1: Read GoldenGate's trail files directly using Python libraries that parse the trail file format (you'll need to handle the binary format specifics).
  • Option 2: Use GoldenGate's REST API (available in newer versions) to fetch change data programmatically.

This is best for large-scale, high-throughput CDC scenarios, but requires setting up and configuring GoldenGate first.

Key Considerations

  • Permissions: Ensure your Python database user has the necessary privileges for the method you choose (trigger creation, LogMiner access, etc.).
  • Performance: Trigger-based CDC adds overhead to your source table's insert operations; LogMiner is more lightweight as it reads logs passively.
  • Data Deduplication: Always track a unique identifier (like changelog_id or SCN) to avoid processing the same change multiple times.
  • Error Handling: Add try/except blocks in your Python code to handle database connection issues or query errors gracefully.

Hope these approaches work for your use case—feel free to tweak them based on your specific data volume and performance needs!

内容的提问来源于stack exchange,提问作者Jeeva Bharathi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:50:39