Python 3.6+
Python Driver
Build powerful Python applications with Cassandra!
📦 Installation & Setup
Get started in minutes!
1
Install Python Driver
Using pip (recommended)
# Install latest version
$ pip install cassandra-driver
# Install specific version
$ pip install cassandra-driver==3.29.0
# Install with optional dependencies
$ pip install cassandra-driver[cql]
# Verify installation
$ python -c "import cassandra; print(cassandra.__version__)"
3.29.0
2
System Requirements
| Component | Requirement | Notes |
|---|---|---|
| Python | 3.6 - 3.12 | Python 2.7 deprecated |
| Cassandra | 2.1+ | Best with 3.11 or 4.0+ |
| Dependencies | six, geomet (optional) | Auto-installed with pip |
| C Extension | libev (optional) | For async performance boost |
Virtual Environment Recommended
# Create virtual environment
$ python -m venv cassandra_env
# Activate (Linux/Mac)
$ source cassandra_env/bin/activate
# Activate (Windows)
$ cassandra_env\Scripts\activate
# Install driver in venv
$ pip install cassandra-driver
🚀 Quick Start Example
Your first Python + Cassandra app!
# hello_cassandra.py
from cassandra.cluster import Cluster
# 1. Connect to Cassandra
cluster = Cluster(['127.0.0.1'])
session = cluster.connect()
# 2. Create keyspace
session.execute("""
CREATE KEYSPACE IF NOT EXISTS tutorial
WITH replication = {
'class': 'SimpleStrategy',
'replication_factor': 1
}
""")
# 3. Use keyspace
session.set_keyspace('tutorial')
# 4. Create table
session.execute("""
CREATE TABLE IF NOT EXISTS users (
user_id UUID PRIMARY KEY,
name TEXT,
email TEXT,
age INT
)
""")
# 5. Insert data
from uuid import uuid4
user_id = uuid4()
session.execute("""
INSERT INTO users (user_id, name, email, age)
VALUES (%s, %s, %s, %s)
""", (user_id, 'Alice', 'alice@example.com', 28))
# 6. Query data
rows = session.execute("SELECT * FROM users")
for row in rows:
print(f"{row.name} ({row.email}) - Age: {row.age}")
# 7. Close connection
cluster.shutdown()
# Output: Alice (alice@example.com) - Age: 28
Run the Example
$ python hello_cassandra.py
Alice (alice@example.com) - Age: 28
🔌 Connection Management
Proper connection setup!
Basic Connection
from cassandra.cluster import Cluster
# Simple connection
cluster = Cluster(['node1', 'node2', 'node3'])
session = cluster.connect('my_keyspace')
# Use session for queries
result = session.execute("SELECT * FROM users")
# Always close when done
cluster.shutdown()
Production-Ready Connection
from cassandra.cluster import Cluster, ExecutionProfile
from cassandra.policies import (
DCAwareRoundRobinPolicy,
TokenAwarePolicy,
DowngradingConsistencyRetryPolicy
)
from cassandra.auth import PlainTextAuthProvider
# Configure authentication
auth_provider = PlainTextAuthProvider(
username='cassandra',
password='cassandra'
)
# Configure load balancing
load_balancing_policy = TokenAwarePolicy(
DCAwareRoundRobinPolicy(local_dc='datacenter1')
)
# Create execution profile
profile = ExecutionProfile(
load_balancing_policy=load_balancing_policy,
retry_policy=DowngradingConsistencyRetryPolicy(),
request_timeout=15.0 # 15 seconds
)
# Create cluster with all options
cluster = Cluster(
contact_points=['node1', 'node2', 'node3'],
port=9042,
auth_provider=auth_provider,
execution_profiles={'default': profile},
protocol_version=5 # Use latest
)
# Connect
session = cluster.connect('my_keyspace')
Singleton Pattern (Recommended)
# cassandra_connection.py
from cassandra.cluster import Cluster
class CassandraConnection:
"""Singleton connection manager"""
_cluster = None
_session = None
@classmethod
def get_session(cls, keyspace='my_keyspace'):
"""Get or create session"""
if cls._session is None:
cls._cluster = Cluster(['node1', 'node2'])
cls._session = cls._cluster.connect(keyspace)
return cls._session
@classmethod
def shutdown(cls):
"""Close connections"""
if cls._cluster:
cls._cluster.shutdown()
cls._cluster = None
cls._session = None
# Usage in your app
from cassandra_connection import CassandraConnection
# In request handlers
def get_user(user_id):
session = CassandraConnection.get_session()
result = session.execute(
"SELECT * FROM users WHERE user_id = %s",
(user_id,)
)
return result.one()
# On app shutdown
CassandraConnection.shutdown()
📝 Basic Query Operations
CRUD operations!
SELECT Queries
# Simple select
rows = session.execute("SELECT * FROM users")
for row in rows:
print(row.name, row.email)
# With WHERE clause
rows = session.execute("""
SELECT name, email FROM users
WHERE user_id = %s
""", (user_id,))
# Get single row
row = session.execute(
"SELECT * FROM users WHERE user_id = %s",
(user_id,)
).one()
if row:
print(f"Name: {row.name}")
# Get first result or None
row = session.execute(
"SELECT * FROM users WHERE email = %s",
(email,)
).one_or_none()
if row is None:
print("User not found")
INSERT Queries
from uuid import uuid4
from datetime import datetime
# Simple insert
session.execute("""
INSERT INTO users (user_id, name, email, age)
VALUES (%s, %s, %s, %s)
""", (uuid4(), 'Bob', 'bob@example.com', 30))
# Insert with timestamp
session.execute("""
INSERT INTO events (event_id, user_id, event_type, timestamp)
VALUES (%s, %s, %s, %s)
""", (uuid4(), user_id, 'login', datetime.now()))
# Insert if not exists
result = session.execute("""
INSERT INTO users (user_id, name, email)
VALUES (%s, %s, %s)
IF NOT EXISTS
""", (user_id, 'Charlie', 'charlie@example.com'))
if result.was_applied:
print("User created")
else:
print("User already exists")
UPDATE Queries
# Simple update
session.execute("""
UPDATE users
SET age = %s, email = %s
WHERE user_id = %s
""", (31, 'newemail@example.com', user_id))
# Conditional update
result = session.execute("""
UPDATE users
SET age = %s
WHERE user_id = %s
IF age = %s
""", (32, user_id, 31))
if result.was_applied:
print("Update successful")
# Increment counter
session.execute("""
UPDATE page_views
SET views = views + 1
WHERE page_id = %s
""", (page_id,))
DELETE Queries
# Delete row
session.execute("""
DELETE FROM users
WHERE user_id = %s
""", (user_id,))
# Delete specific columns
session.execute("""
DELETE email, age FROM users
WHERE user_id = %s
""", (user_id,))
# Conditional delete
result = session.execute("""
DELETE FROM users
WHERE user_id = %s
IF age > %s
""", (user_id, 50))
if result.was_applied:
print("User deleted")
BATCH Queries
from cassandra.query import BatchStatement
# Create batch
batch = BatchStatement()
# Add statements to batch
batch.add("INSERT INTO users (user_id, name) VALUES (%s, %s)",
(uuid4(), 'User1'))
batch.add("INSERT INTO users (user_id, name) VALUES (%s, %s)",
(uuid4(), 'User2'))
batch.add("INSERT INTO users (user_id, name) VALUES (%s, %s)",
(uuid4(), 'User3'))
# Execute batch (atomic)
session.execute(batch)
# WARNING: Only batch operations on SAME partition!
# Don't batch unrelated data!
⚡ Prepared Statements
10x faster queries!
Why Use Prepared Statements?
- ✅ Performance: Query parsed once, executed many times (10x faster!)
- ✅ Security: Protection from CQL injection
- ✅ Type checking: Parameters validated
- ✅ Network: Less data sent to Cassandra
Basic Prepared Statements
# Prepare statement ONCE (at startup)
insert_user = session.prepare("""
INSERT INTO users (user_id, name, email, age)
VALUES (?, ?, ?, ?)
""")
# Execute MANY times (in request handlers)
session.execute(insert_user, (uuid4(), 'Alice', 'alice@example.com', 28))
session.execute(insert_user, (uuid4(), 'Bob', 'bob@example.com', 30))
session.execute(insert_user, (uuid4(), 'Charlie', 'charlie@example.com', 25))
# 10x faster than simple queries! ✅
Named Parameters
# Prepare with named parameters
select_user = session.prepare("""
SELECT name, email, age FROM users
WHERE user_id = :user_id AND age > :min_age
""")
# Execute with dict
rows = session.execute(
select_user,
{'user_id': user_id, 'min_age': 25}
)
for row in rows:
print(f"{row.name}: {row.age} years old")
Prepared Statement Pattern
# prepared_statements.py
class UserQueries:
"""Pre-prepared query statements"""
def __init__(self, session):
# Prepare all statements at startup
self.insert_user = session.prepare("""
INSERT INTO users (user_id, name, email, age)
VALUES (?, ?, ?, ?)
""")
self.get_user = session.prepare("""
SELECT * FROM users WHERE user_id = ?
""")
self.update_user = session.prepare("""
UPDATE users SET name = ?, email = ?, age = ?
WHERE user_id = ?
""")
self.delete_user = session.prepare("""
DELETE FROM users WHERE user_id = ?
""")
def create_user(self, session, user_id, name, email, age):
session.execute(self.insert_user, (user_id, name, email, age))
def find_user(self, session, user_id):
return session.execute(self.get_user, (user_id,)).one()
# Usage
queries = UserQueries(session)
queries.create_user(session, uuid4(), 'Alice', 'alice@example.com', 28)
🔄 Asynchronous Queries
Non-blocking operations!
execute_async()
# Execute query asynchronously
future = session.execute_async("""
SELECT * FROM users WHERE user_id = %s
""", (user_id,))
# Do other work while query executes...
print("Query running in background...")
# Get result (blocks until ready)
try:
rows = future.result()
for row in rows:
print(row.name)
except Exception as e:
print(f"Query failed: {e}")
Callbacks
# Define success callback
def handle_success(rows):
print("Query succeeded!")
for row in rows:
print(f"User: {row.name}")
# Define error callback
def handle_error(exception):
print(f"Query failed: {exception}")
# Execute with callbacks
future = session.execute_async("SELECT * FROM users")
future.add_callback(handle_success)
future.add_errback(handle_error)
# Continue with other work
# Callbacks execute when query completes
Parallel Queries
# Execute multiple queries in parallel
user_ids = [uuid1, uuid2, uuid3, uuid4, uuid5]
# Start all queries
futures = []
for user_id in user_ids:
future = session.execute_async(
"SELECT * FROM users WHERE user_id = %s",
(user_id,)
)
futures.append(future)
# Wait for all to complete
results = []
for future in futures:
try:
row = future.result().one()
results.append(row)
except Exception as e:
print(f"Query failed: {e}")
# Process all results
for row in results:
print(f"Found user: {row.name}")
# Much faster than sequential queries! ⚡
🚨 Error Handling
Handle failures gracefully!
Common Exceptions
from cassandra import (
ReadTimeout,
WriteTimeout,
Unavailable,
InvalidRequest,
NoHostAvailable
)
try:
session.execute("SELECT * FROM users WHERE user_id = %s", (user_id,))
except ReadTimeout:
print("Query timed out - try again")
except WriteTimeout:
print("Write timed out - data may or may not be written")
except Unavailable as e:
print(f"Not enough replicas: needed {e.required_replicas}, alive {e.alive_replicas}")
except InvalidRequest as e:
print(f"Bad query: {e}")
except NoHostAvailable:
print("Cannot connect to cluster")
except Exception as e:
print(f"Unexpected error: {e}")
Retry Logic
import time
def execute_with_retry(session, query, params, max_retries=3):
"""Execute query with exponential backoff retry"""
for attempt in range(max_retries):
try:
return session.execute(query, params)
except (ReadTimeout, WriteTimeout) as e:
if attempt == max_retries - 1:
raise # Last attempt failed
# Exponential backoff: 1s, 2s, 4s
wait_time = 2 ** attempt
print(f"Retry {attempt + 1}/{max_retries} after {wait_time}s")
time.sleep(wait_time)
raise Exception("All retries failed")
# Usage
result = execute_with_retry(
session,
"SELECT * FROM users WHERE user_id = %s",
(user_id,)
)
💡 Python Driver Best Practices
Production tips!
DO
- Use singleton pattern for cluster/session
- Use prepared statements for repeated queries
- Use async queries for parallelism
- Handle exceptions properly
- Use parameterized queries (%s or ?)
- Close cluster on shutdown
- Use connection pooling (automatic)
- Monitor session metrics
DON'T
- Create cluster/session per request
- Use string concatenation for queries
- Ignore errors silently
- Use execute() for repeated queries
- Forget to close cluster
- Use SELECT * on large tables
- Batch unrelated operations
- Skip query timeouts
Performance Tips
| Technique | Benefit | When to Use |
|---|---|---|
| Prepared Statements | 10x faster execution | Queries executed multiple times |
| Async Queries | Non-blocking, parallel | Multiple independent queries |
| Paging | Lower memory usage | Large result sets |
| Token-aware routing | Fewer network hops | Always (enable in policy) |
| Compression | Reduced bandwidth | High-latency networks |
Complete Example
# user_repository.py - Production-ready pattern
from cassandra.cluster import Cluster
from cassandra.query import dict_factory
from uuid import uuid4
class UserRepository:
"""User data access layer"""
_cluster = None
_session = None
_queries = {}
@classmethod
def initialize(cls, contact_points):
"""Initialize connection (call at startup)"""
cls._cluster = Cluster(contact_points)
cls._session = cls._cluster.connect('my_keyspace')
cls._session.row_factory = dict_factory # Return dicts
# Prepare statements
cls._queries['insert'] = cls._session.prepare("""
INSERT INTO users (user_id, name, email, age)
VALUES (?, ?, ?, ?)
""")
cls._queries['get'] = cls._session.prepare("""
SELECT * FROM users WHERE user_id = ?
""")
@classmethod
def create_user(cls, name, email, age):
"""Create new user"""
user_id = uuid4()
cls._session.execute(
cls._queries['insert'],
(user_id, name, email, age)
)
return user_id
@classmethod
def get_user(cls, user_id):
"""Get user by ID"""
result = cls._session.execute(
cls._queries['get'],
(user_id,)
)
return result.one()
@classmethod
def shutdown(cls):
"""Close connections (call on shutdown)"""
if cls._cluster:
cls._cluster.shutdown()
# Usage in your app
if __name__ == '__main__':
# Startup
UserRepository.initialize(['localhost'])
# Use in request handlers
user_id = UserRepository.create_user('Alice', 'alice@example.com', 28)
user = UserRepository.get_user(user_id)
print(user)
# Shutdown
UserRepository.shutdown()
🎉 Master Python + Cassandra!
You now know how to build production Python apps with Cassandra!
🎓 What You Learned:
- 📦 Installation: pip install cassandra-driver (Python 3.6+)
- 🚀 Quick start: Connect, create tables, CRUD in 20 lines
- 🔌 Connection: Singleton pattern for proper management
- 📝 Basic queries: SELECT, INSERT, UPDATE, DELETE, BATCH
- ⚡ Prepared statements: 10x faster with prepare()
- 🔄 Async queries: execute_async() for parallelism
- 🚨 Error handling: Try/except with proper retry logic
- 💡 Best practices: Production patterns and tips
💡 Key Takeaways:
- Singleton pattern - One cluster/session for entire app
- Prepared statements - Use for repeated queries (10x faster!)
- Async queries - execute_async() for parallel operations
- Error handling - Catch specific exceptions, implement retry
- Row factory - dict_factory for dictionary results
- Shutdown properly - cluster.shutdown() on exit
🐍 Quick Reference:
# Install
pip install cassandra-driver
# Connect (ONCE at startup)
from cassandra.cluster import Cluster
cluster = Cluster(['node1', 'node2'])
session = cluster.connect('keyspace')
# Prepared statement (10x faster!)
stmt = session.prepare("SELECT * FROM users WHERE id = ?")
result = session.execute(stmt, (user_id,))
# Async query (parallel!)
future = session.execute_async("SELECT ...")
rows = future.result()
# Shutdown
cluster.shutdown()
🐍 Python + Cassandra = Powerful combo! 🎯
Start building your app now!
Advertisement
📱 Responsive Ad 📱