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")