Step 5: Create the event consumer script

Teradata Developer Guides

ft:locale
en-US
ft:lastEdition
2026-08-18

Create a file named consumer.py:

import redis
import teradatasql
import os
from datetime import datetime
from dotenv import load_dotenv

# Load environment variables from .env file
load_dotenv()

# Redis connection
redis_host = os.getenv("REDIS_HOST", "localhost")
redis_port = int(os.getenv("REDIS_PORT", 6379))

# Teradata connection
teradata_host = os.getenv("TERADATA_HOST")
teradata_user = os.getenv("TERADATA_USER")
teradata_password = os.getenv("TERADATA_PASSWORD")

if not all([teradata_host, teradata_user, teradata_password]):
    raise ValueError("Missing TERADATA_HOST, TERADATA_USER, or TERADATA_PASSWORD environment variables")

# Connect to Redis
try:
    r = redis.Redis(host=redis_host, port=redis_port, decode_responses=True)
    r.ping()
    print(f"✓ Connected to Redis at {redis_host}:{redis_port}")
except Exception as e:
    print(f"✗ Failed to connect to Redis at {redis_host}:{redis_port}")
    print(f"  Error: {e}")
    exit(1)

# Connect to Teradata
try:
    conn = teradatasql.connect(
        host=teradata_host,
        user=teradata_user,
        password=teradata_password
    )
    print(f"✓ Connected to Teradata at {teradata_host}")
except Exception as e:
    print(f"✗ Failed to connect to Teradata")
    print(f"  Error: {e}")
    exit(1)

cursor = conn.cursor()

STREAM_KEY = "events:stream"
CONSUMER_GROUP = "teradata-consumer"
CONSUMER_NAME = "consumer-1"

print()
print("Redis Event Consumer")
print("=" * 50)
print(f"Listening on stream: {STREAM_KEY}")
print()

# Create consumer group if it doesn't exist
try:
    r.xgroup_create(STREAM_KEY, CONSUMER_GROUP, id="0", mkstream=True)
    print("✓ Consumer group created")
except redis.ResponseError as e:
    if "already exists" in str(e):
        print("✓ Consumer group already exists")
    else:
        raise

print()

# Read events from the stream
events_processed = 0

try:
    while True:
        # Read from consumer group
        messages = r.xreadgroup(
            groupname=CONSUMER_GROUP,
            consumername=CONSUMER_NAME,
            streams={STREAM_KEY: ">"},
            count=10,
            block=1000  # Block for 1 second if no messages
        )

        if not messages:
            print("Waiting for events...")
            continue

        for stream_key, message_list in messages:
            for message_id, message_data in message_list:
                try:
                    # Extract event data
                    event_id = int(message_data.get("event_id", 0))
                    event_type = message_data.get("event_type", "unknown")
                    user_id = int(message_data.get("user_id", 0))
                    event_timestamp = message_data.get("event_timestamp", datetime.now().strftime("%Y-%m-%d %H:%M:%S"))
                    country = message_data.get("country", "Unknown")

                    # Insert into Teradata using parameterized query
                    insert_query = """
                    INSERT INTO event_data (
                        event_id, event_type, user_id, event_timestamp, country
                    ) VALUES (
                        ?, ?, ?, ?, ?
                    );
                    """

                    cursor.execute(insert_query, (event_id, event_type, user_id, event_timestamp, country))
                    conn.commit()

                    # Acknowledge the message only after successful commit
                    r.xack(STREAM_KEY, CONSUMER_GROUP, message_id)

                    events_processed += 1
                    print(f"✓ Event {events_processed}: {event_type} (user={user_id}) → Teradata")

                except Exception as e:
                    conn.rollback()
                    print(f"✗ Error processing event: {e}")

except KeyboardInterrupt:
    print()
    print("✓ Consumer stopped")
finally:
    cursor.close()
    conn.close()
    r.close()
    print(f"✓ Processed {events_processed} events total")