Real-time data analysis has become a critical component for companies that need to make fast decisions based on up-to-date information. Python has established itself as one of the most powerful languages for implementing real-time data processing solutions, combining ease of use with robust libraries. This article explores how to use Python for real-time analysis, focusing on the main technologies and techniques.
Python has consolidated its position as the preferred language for real-time data analysis for several reasons:
These advantages are becoming even more evident as the Python ecosystem for real-time analysis evolves to meet growing demands for data speed and volume.
Before we dive into specific tools, it is important to understand the common architectures for real-time data analysis with Python:
The Lambda architecture combines batch and real-time processing. It is useful when I need both immediate and historical data for analysis:
[Data Sources] → [Speed Layer (Real Time)] → [Serving Layer] → [Applications] → [Batch Layer] →
The Kappa architecture simplifies the model by treating everything as streams. It is ideal when the focus is entirely on real-time data and continuous events:
[Data Sources] → [Streaming System] → [Stream Processing] → [Storage] → [Applications]
The SMACK stack (Spark, Mesos, Akka, Cassandra, Kafka) has become popular for real-time data applications:
[Kafka (Ingestion)] → [Spark Streaming (Processing)] → [Cassandra (Storage)] → [Applications]
↑
[Mesos/Kubernetes (Orchestration)]
↑
[Akka (Messaging)]
The Python ecosystem for real-time analysis offers several powerful tools. Let’s explore the most important ones.
Apache Kafka remains the most popular distributed streaming platform, and its integration with Python is excellent.
Apache Kafka is one of the top options when we need a high-performance distributed message queue. We can use the confluent-kafka library to integrate it with Python.
Before showing the code, it is important to understand that I am simulating a scenario of sensors sending temperature and humidity readings to a Kafka topic. The producer sends these readings, while the consumer processes them and identifies, for example, high-temperature alerts.
# Kafka producer example with confluent-kafka
from confluent_kafka import Producer
import json
import time
import random
# Producer configuration
config = {
'bootstrap.servers': 'localhost:9092',
'client.id': 'python-producer'
}
producer = Producer(config)
# Callback function to confirm delivery
def delivery_report(err, msg):
if err is not None:
print(f'Message delivery failed: {err}')
else:
print(f'Message delivered to topic {msg.topic()} [partition {msg.partition()}]')
# Generate and send simulated data
for i in range(100):
# Create simulated sensor data
data = {
'sensor_id': f'sensor-{random.randint(1, 10)}',
'temperature': round(random.uniform(20.0, 35.0), 2),
'humidity': round(random.uniform(30.0, 80.0), 2),
'timestamp': int(time.time())
}
# Serialize to JSON
payload = json.dumps(data)
# Send to the Kafka topic
producer.produce('sensor-data',
key=data['sensor_id'],
value=payload,
callback=delivery_report)
# Flush the buffer periodically
producer.poll(0)
time.sleep(0.5) # Simulate the interval between readings
# Make sure all messages are sent
producer.flush()
This consumer reads the data that was sent, converts it from JSON, and checks whether the temperature is above a predefined threshold. It is a simple foundation for anomaly detection logic.
# Kafka consumer example with confluent-kafka
from confluent_kafka import Consumer, KafkaError
import json
# Consumer configuration
config = {
'bootstrap.servers': 'localhost:9092',
'group.id': 'python-consumer-group',
'auto.offset.reset': 'earliest'
}
consumer = Consumer(config)
consumer.subscribe(['sensor-data'])
try:
while True:
msg = consumer.poll(1.0)
if msg is None:
continue
if msg.error():
if msg.error().code() == KafkaError._PARTITION_EOF:
print(f'Reached the end of partition {msg.partition()}')
else:
print(f'Error: {msg.error()}')
else:
# Process the received message
try:
data = json.loads(msg.value())
print(f"Sensor: {data['sensor_id']}, Temperature: {data['temperature']}°C, Humidity: {data['humidity']}%")
# This is where you could implement real-time processing logic
if data['temperature'] > 30.0:
print(f"ALERT: High temperature detected on sensor {data['sensor_id']}!")
except json.JSONDecodeError:
print("Error decoding JSON")
except KeyboardInterrupt:
pass
finally:
# Close the consumer
consumer.close()
Apache Spark, with its Python API (PySpark), offers powerful stream processing capabilities.
When we need more aggregation power and time-window analysis, Spark Structured Streaming with PySpark becomes an essential tool.
In the example below, I connect a Kafka stream to Spark to calculate per-sensor statistics over 5-minute sliding windows and detect anomalies.
# Spark Structured Streaming example with PySpark
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark.sql.types import *
# Create a Spark session
spark = SparkSession.builder \
.appName("RealTimeAnalytics") \
.config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.4.0") \
.getOrCreate()
# Define the data schema
schema = StructType([
StructField("sensor_id", StringType(), True),
StructField("temperature", FloatType(), True),
StructField("humidity", FloatType(), True),
StructField("timestamp", TimestampType(), True)
])
# Read the stream from Kafka
kafka_stream = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "sensor-data") \
.option("startingOffsets", "latest") \
.load()
# Extract and transform the data
parsed_stream = kafka_stream \
.select(from_json(col("value").cast("string"), schema).alias("data")) \
.select("data.*") \
.withWatermark("timestamp", "10 seconds")
# Calculate statistics over time windows
window_stats = parsed_stream \
.groupBy(
window(col("timestamp"), "5 minutes", "1 minute"),
col("sensor_id")
) \
.agg(
avg("temperature").alias("avg_temp"),
max("temperature").alias("max_temp"),
min("temperature").alias("min_temp"),
avg("humidity").alias("avg_humidity")
)
# Detect anomalies
anomalies = parsed_stream \
.filter(col("temperature") > 32.0) \
.select(
col("sensor_id"),
col("temperature"),
col("timestamp")
)
# Console output (for development)
query1 = window_stats.writeStream \
.outputMode("complete") \
.format("console") \
.option("truncate", "false") \
.start()
query2 = anomalies.writeStream \
.outputMode("append") \
.format("console") \
.start()
# Database output (for production)
# PostgreSQL example
postgres_output = window_stats.writeStream \
.foreachBatch(lambda df, epoch_id: df.write \
.format("jdbc") \
.option("url", "jdbc:postgresql://localhost:5432/sensordb") \
.option("dbtable", "sensor_stats") \
.option("user", "username") \
.option("password", "password") \
.mode("append") \
.save()
) \
.outputMode("update") \
.start()
# Wait for termination (in production, this would be controlled externally)
spark.streams.awaitAnyTermination()
Apache Flink has gained popularity for its low-latency stream processing.
Apache Flink is a great choice when we need extremely low latency and fine-grained control over stream events. PyFlink offers a powerful API based on SQL and the DataStream API.
Before showing the code, note that here I configure a Kafka source, process events per minute with aggregations, and send the results back to another Kafka topic.
# PyFlink example for stream processing
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment, EnvironmentSettings
from pyflink.table.descriptors import Schema, Kafka, Json
from pyflink.table.window import Tumble
# Create the execution environment
env = StreamExecutionEnvironment.get_execution_environment()
env_settings = EnvironmentSettings.new_instance().in_streaming_mode().build()
t_env = StreamTableEnvironment.create(env, environment_settings=env_settings)
# Configure the Kafka connection
t_env.connect(
Kafka()
.version("universal")
.topic("sensor-data")
.start_from_latest()
.property("bootstrap.servers", "localhost:9092")
.property("group.id", "pyflink-consumer")
) \
.with_format(
Json()
.fail_on_missing_field(False)
.schema(
Schema()
.field("sensor_id", "STRING")
.field("temperature", "DOUBLE")
.field("humidity", "DOUBLE")
.field("timestamp", "BIGINT")
)
) \
.with_schema(
Schema()
.field("sensor_id", "STRING")
.field("temperature", "DOUBLE")
.field("humidity", "DOUBLE")
.field("event_time", "TIMESTAMP(3)")
.proctime()
) \
.create_temporary_table("sensor_source")
# Define the SQL query for real-time processing
result_table = t_env.sql_query("""
SELECT
sensor_id,
TUMBLE_START(event_time, INTERVAL '1' MINUTE) AS window_start,
TUMBLE_END(event_time, INTERVAL '1' MINUTE) AS window_end,
COUNT(*) AS reading_count,
AVG(temperature) AS avg_temp,
MAX(temperature) AS max_temp,
MIN(temperature) AS min_temp,
AVG(humidity) AS avg_humidity
FROM sensor_source
GROUP BY
TUMBLE(event_time, INTERVAL '1' MINUTE),
sensor_id
""")
# Configure console output (development)
t_env.connect(
Kafka()
.version("universal")
.topic("sensor-stats")
.property("bootstrap.servers", "localhost:9092")
.start_from_latest()
) \
.with_format(
Json()
.derive_schema()
) \
.with_schema(
Schema()
.field("sensor_id", "STRING")
.field("window_start", "TIMESTAMP(3)")
.field("window_end", "TIMESTAMP(3)")
.field("reading_count", "BIGINT")
.field("avg_temp", "DOUBLE")
.field("max_temp", "DOUBLE")
.field("min_temp", "DOUBLE")
.field("avg_humidity", "DOUBLE")
) \
.create_temporary_table("sensor_sink")
# Run the insert into the output table
result_table.execute_insert("sensor_sink").wait()
For simpler use cases, native Python libraries offer efficient solutions.
When we need a lightweight, fast solution that is easy to prototype, the streamz library helps us build streaming pipelines directly in pure Python. It is ideal for local testing or smaller applications, especially when I am not using distributed clusters.
In the following example, I simulate a stream of sensor data and apply parsing, filtering, and basic statistical analysis with Pandas:
# Example with streamz for stream processing in pure Python
from streamz import Stream
import json
import time
import random
import pandas as pd
from datetime import datetime
# Create the stream
source = Stream()
# Define the processing functions
def parse_data(message):
try:
return json.loads(message)
except json.JSONDecodeError:
return None
def filter_valid_readings(data):
if data is None:
return False
return 'sensor_id' in data and 'temperature' in data and 'humidity' in data
def detect_anomalies(data):
if data['temperature'] > 30.0:
print(f"ALERT: High temperature ({data['temperature']}°C) detected on sensor {data['sensor_id']}!")
return data
def add_timestamp(data):
data['processed_at'] = datetime.now().isoformat()
return data
def batch_to_dataframe(batch):
df = pd.DataFrame(batch)
return df
# Build the processing pipeline
processed = source \
.map(parse_data) \
.filter(filter_valid_readings) \
.map(detect_anomalies) \
.map(add_timestamp)
# Create sliding windows for analysis
windowed = processed \
.sliding_window(10) \
.map(batch_to_dataframe) \
.map(lambda df: df.describe())
# Add outputs
processed.sink(print)
windowed.sink(print)
# Simulate a data source
for i in range(100):
data = {
'sensor_id': f'sensor-{random.randint(1, 5)}',
'temperature': round(random.uniform(20.0, 35.0), 2),
'humidity': round(random.uniform(30.0, 80.0), 2),
'timestamp': int(time.time())
}
source.emit(json.dumps(data))
time.sleep(0.5)
Real-time data analysis with Python has applications across many industries. Let’s explore some popular use cases:
One interesting use case is IoT projects. Here, sensors spread across a site send data to a pipeline that detects unusual patterns. We use Scikit-learn to apply simple models such as Isolation Forest.
Before showing the code, it is worth understanding that the model is trained on historical data and then applied in real time as readings arrive through Kafka.
# Real-time anomaly detection example with Scikit-learn and Kafka
from confluent_kafka import Consumer
import json
import numpy as np
from sklearn.ensemble import IsolationForest
import pandas as pd
import time
# Configure the anomaly detection model
model = IsolationForest(contamination=0.05, random_state=42)
# Data history for initial training
historical_data = []
# Configure the Kafka consumer
consumer = Consumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'anomaly-detector',
'auto.offset.reset': 'earliest'
})
consumer.subscribe(['sensor-data'])
# Function to retrain the model
def retrain_model(data_points):
if len(data_points) < 100:
return False
# Convert to a DataFrame
df = pd.DataFrame(data_points)
# Select only numeric columns for training
features = df[['temperature', 'humidity']].values
# Train the model
model.fit(features)
print(f"Model retrained with {len(data_points)} data points")
return True
# Function to detect anomalies
def detect_anomaly(data_point):
# Extract features
features = np.array([[data_point['temperature'], data_point['humidity']]])
# Predict
prediction = model.predict(features)
score = model.decision_function(features)
# -1 means anomaly, 1 means normal
is_anomaly = prediction[0] == -1
return is_anomaly, score[0]
# Process the stream
try:
# Initial phase: collect data for training
print("Collecting initial data for training...")
start_time = time.time()
while len(historical_data) < 100 and time.time() - start_time < 60:
msg = consumer.poll(1.0)
if msg is None:
continue
if msg.error():
print(f"Error: {msg.error()}")
continue
try:
data = json.loads(msg.value())
historical_data.append(data)
print(f"Collected: {len(historical_data)}/100 data points")
except json.JSONDecodeError:
print("Error decoding JSON")
# Train the initial model
if len(historical_data) >= 50:
retrain_model(historical_data)
print("Initial model trained. Starting anomaly detection...")
else:
print("Not enough data for initial training.")
exit(1)
# Detection phase
retrain_counter = 0
while True:
msg = consumer.poll(1.0)
if msg is None:
continue
if msg.error():
print(f"Error: {msg.error()}")
continue
try:
data = json.loads(msg.value())
# Detect anomaly
is_anomaly, score = detect_anomaly(data)
# Add to history
historical_data.append(data)
if len(historical_data) > 1000: # Keep a sliding window
historical_data.pop(0)
# Report the result
if is_anomaly:
print(f"ANOMALY DETECTED! Sensor: {data['sensor_id']}, Temp: {data['temperature']}°C, Humidity: {data['humidity']}%, Score: {score:.4f}")
else:
print(f"Normal - Sensor: {data['sensor_id']}, Score: {score:.4f}")
# Retrain periodically
retrain_counter += 1
if retrain_counter >= 100:
retrain_model(historical_data)
retrain_counter = 0
except json.JSONDecodeError:
print("Error decoding JSON")
except KeyboardInterrupt:
pass
finally:
consumer.close()
Another very interesting case where we can apply Python for real-time data analysis is social media monitoring. Here, the goal is to classify the sentiment of posts as they are published, using pre-trained NLP models such as distilbert-base-uncased-finetuned-sst-2-english.
The idea is to consume posts from a Kafka topic, apply the sentiment analysis model with the transformers library, and send the results to another topic for visualization or storage.
# Real-time sentiment analysis example for tweets
from confluent_kafka import Consumer, Producer
import json
from transformers import pipeline
import time
# Initialize the sentiment analysis model
sentiment_analyzer = pipeline("sentiment-analysis", model="distilbert-base-uncased-finetuned-sst-2-english")
# Configure the Kafka consumer
consumer = Consumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'sentiment-analyzer',
'auto.offset.reset': 'earliest'
})
consumer.subscribe(['social-media-posts'])
# Configure the Kafka producer for results
producer = Producer({
'bootstrap.servers': 'localhost:9092'
})
# Function to analyze sentiment
def analyze_sentiment(text):
try:
result = sentiment_analyzer(text)[0]
return {
'label': result['label'],
'score': float(result['score'])
}
except Exception as e:
print(f"Sentiment analysis error: {e}")
return {
'label': 'ERROR',
'score': 0.0
}
# Process the stream
try:
while True:
msg = consumer.poll(1.0)
if msg is None:
continue
if msg.error():
print(f"Error: {msg.error()}")
continue
try:
# Decode the message
post = json.loads(msg.value())
# Extract the text
text = post.get('text', '')
if not text:
continue
# Analyze sentiment
sentiment = analyze_sentiment(text)
# Add the result to the post
post['sentiment'] = sentiment
# Send the result to another topic
producer.produce(
'analyzed-posts',
key=post.get('id', str(time.time())),
value=json.dumps(post)
)
# Print the result
print(f"Post: '{text[:50]}...' - Sentiment: {sentiment['label']} ({sentiment['score']:.4f})")
# Flush the buffer periodically
producer.poll(0)
except json.JSONDecodeError:
print("Error decoding JSON")
except KeyboardInterrupt:
pass
finally:
consumer.close()
producer.flush()
In many projects, especially in a corporate context, we need to visualize data processed in real time in a clear, interactive way. For that, we use the Dash framework together with Plotly, creating dashboards that update automatically with the latest readings received through Kafka.
The example below shows how I build a real-time dashboard to visualize sensor temperature and humidity, consuming data from a Kafka topic and displaying charts that update every second.
# Real-time dashboard example with Dash and Kafka
import dash
from dash import dcc, html
from dash.dependencies import Input, Output
import plotly.graph_objs as go
from confluent_kafka import Consumer
import json
import pandas as pd
from collections import deque
import threading
import time
# Configure in-memory data storage
max_length = 100
times = deque(maxlen=max_length)
temperatures = {f'sensor-{i}': deque(maxlen=max_length) for i in range(1, 6)}
humidities = {f'sensor-{i}': deque(maxlen=max_length) for i in range(1, 6)}
# Function that consumes Kafka data in a separate thread
def consume_kafka_data():
# Configure the consumer
consumer = Consumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'dashboard-consumer',
'auto.offset.reset': 'latest'
})
consumer.subscribe(['sensor-data'])
try:
while True:
msg = consumer.poll(1.0)
if msg is None:
continue
if msg.error():
print(f"Error: {msg.error()}")
continue
try:
# Process the message
data = json.loads(msg.value())
# Extract the data
sensor_id = data.get('sensor_id')
temperature = data.get('temperature')
humidity = data.get('humidity')
timestamp = data.get('timestamp')
# Add to the deques
if sensor_id in temperatures and temperature is not None:
current_time = pd.to_datetime(timestamp, unit='s')
times.append(current_time)
temperatures[sensor_id].append(temperature)
humidities[sensor_id].append(humidity)
except json.JSONDecodeError:
print("Error decoding JSON")
except Exception as e:
print(f"Kafka consumer error: {e}")
finally:
consumer.close()
# Start the consumer thread
kafka_thread = threading.Thread(target=consume_kafka_data, daemon=True)
kafka_thread.start()
# Create the Dash application
app = dash.Dash(__name__)
app.layout = html.Div([
html.H1("Real-Time Sensor Dashboard"),
html.Div([
html.H2("Temperature"),
dcc.Graph(id='temperature-graph'),
dcc.Interval(
id='temperature-update',
interval=1000, # Update every second
n_intervals=0
)
]),
html.Div([
html.H2("Humidity"),
dcc.Graph(id='humidity-graph'),
dcc.Interval(
id='humidity-update',
interval=1000, # Update every second
n_intervals=0
)
])
])
@app.callback(
Output('temperature-graph', 'figure'),
Input('temperature-update', 'n_intervals')
)
def update_temperature_graph(n):
traces = []
for sensor_id, temp_data in temperatures.items():
if len(temp_data) > 0:
traces.append(go.Scatter(
x=list(times)[-len(temp_data):],
y=list(temp_data),
name=sensor_id,
mode='lines+markers'
))
return {
'data': traces,
'layout': go.Layout(
xaxis=dict(title='Time'),
yaxis=dict(title='Temperature (°C)'),
title='Real-Time Sensor Temperature',
height=400
)
}
@app.callback(
Output('humidity-graph', 'figure'),
Input('humidity-update', 'n_intervals')
)
def update_humidity_graph(n):
traces = []
for sensor_id, humidity_data in humidities.items():
if len(humidity_data) > 0:
traces.append(go.Scatter(
x=list(times)[-len(humidity_data):],
y=list(humidity_data),
name=sensor_id,
mode='lines+markers'
))
return {
'data': traces,
'layout': go.Layout(
xaxis=dict(title='Time'),
yaxis=dict(title='Humidity (%)'),
title='Real-Time Sensor Humidity',
height=400
)
}
# Run the server
if __name__ == '__main__':
app.run_server(debug=True, host='0.0.0.0')
This visual approach makes it much easier to track unexpected behavior or trends in the data.
To implement efficient real-time data analysis solutions with Python, consider these best practices:
When processing large volumes of messages, grouping data into batches or time windows reduces overhead and improves performance. Libraries such as Spark and Beam support this kind of aggregation natively.
# Optimization example with batch processing
from confluent_kafka import Consumer
import json
import time
# Configure the consumer
consumer = Consumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'batch-processor',
'auto.offset.reset': 'earliest',
'max.poll.records': 500 # Process up to 500 records at a time
})
consumer.subscribe(['high-volume-data'])
# Process in batches
try:
while True:
batch = []
start_time = time.time()
# Collect a batch by time or size
while time.time() - start_time < 5 and len(batch) < 1000:
msg = consumer.poll(0.1)
if msg is None:
continue
if msg.error():
print(f"Error: {msg.error()}")
continue
try:
data = json.loads(msg.value())
batch.append(data)
except json.JSONDecodeError:
print("Error decoding JSON")
# Process the batch
if batch:
print(f"Processing batch of {len(batch)} messages")
# This is where you would implement batch processing
# For example, using pandas for efficient analysis
# df = pd.DataFrame(batch)
# results = df.groupby('category').agg({'value': ['mean', 'sum', 'count']})
print(f"Batch processed in {time.time() - start_time:.2f} seconds")
except KeyboardInterrupt:
pass
finally:
consumer.close()
Errors happen, especially when consuming external data or working with unstable streams. Implementing patterns such as retries with backoff, circuit breakers, and structured logs keeps things robust.
# Circuit breaker pattern example for external APIs
import requests
import time
from functools import wraps
class CircuitBreaker:
def __init__(self, max_failures=3, reset_timeout=60):
self.max_failures = max_failures
self.reset_timeout = reset_timeout
self.failures = 0
self.state = "CLOSED" # CLOSED, OPEN, HALF-OPEN
self.last_failure_time = None
def __call__(self, func):
@wraps(func)
def wrapper(*args, **kwargs):
if self.state == "OPEN":
# Check whether the reset timeout has passed
if time.time() - self.last_failure_time > self.reset_timeout:
self.state = "HALF-OPEN"
print(f"Circuit breaker switched to HALF-OPEN after {self.reset_timeout}s")
else:
raise Exception(f"Circuit breaker open. Retrying in {self.reset_timeout - (time.time() - self.last_failure_time):.1f}s")
try:
result = func(*args, **kwargs)
# Success in HALF-OPEN state, reset
if self.state == "HALF-OPEN":
self.failures = 0
self.state = "CLOSED"
print("Circuit breaker reset to CLOSED after success")
return result
except Exception as e:
self.failures += 1
self.last_failure_time = time.time()
if self.state == "CLOSED" and self.failures >= self.max_failures:
self.state = "OPEN"
print(f"Circuit breaker opened after {self.failures} failures")
raise e
return wrapper
# Using the circuit breaker
@CircuitBreaker(max_failures=3, reset_timeout=30)
def call_external_api(url):
response = requests.get(url, timeout=5)
response.raise_for_status()
return response.json()
# Example usage in stream processing
def process_stream_with_resilience():
while True:
try:
# Get data from the stream
data = get_next_data_point()
# Call the external API with the circuit breaker
try:
enriched_data = call_external_api(f"https://api.example.com/enrich?id={data['id']}")
data.update(enriched_data)
except Exception as e:
print(f"Error enriching data: {e}")
# Continue with partial data
# Process and save the data
process_and_save(data)
except Exception as e:
print(f"Processing error: {e}")
# Implement exponential backoff
time.sleep(retry_delay)
retry_delay = min(retry_delay * 2, max_retry_delay)
Tools such as Prometheus, Grafana, or OpenTelemetry let us monitor the health of our pipelines. Metrics such as latency, throughput, and errors per second help diagnose problems before they affect the business.
# Instrumentation example with OpenTelemetry
from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.resources import SERVICE_NAME, Resource
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.trace.status import Status, StatusCode
import time
import random
# Configure the tracer
resource = Resource(attributes={
SERVICE_NAME: "real-time-data-processor"
})
provider = TracerProvider(resource=resource)
processor = BatchSpanProcessor(OTLPSpanExporter(endpoint="localhost:4317"))
provider.add_span_processor(processor)
trace.set_tracer_provider(provider)
tracer = trace.get_tracer(__name__)
# Function that processes a message with tracing
def process_message(message_id, payload):
with tracer.start_as_current_span("process_message") as span:
span.set_attribute("message.id", message_id)
span.set_attribute("message.size", len(payload))
try:
# Simulate processing steps
with tracer.start_as_current_span("parse_data"):
time.sleep(random.uniform(0.01, 0.05)) # Simulate work
data = json.loads(payload)
with tracer.start_as_current_span("transform_data"):
time.sleep(random.uniform(0.05, 0.1)) # Simulate work
# Transform the data...
with tracer.start_as_current_span("save_results"):
time.sleep(random.uniform(0.1, 0.2)) # Simulate work
# Save the results...
span.set_status(Status(StatusCode.OK))
return True
except Exception as e:
span.set_status(Status(StatusCode.ERROR, str(e)))
span.record_exception(e)
return False
# Usage in a stream processor
def process_stream_with_tracing():
for i in range(100):
message_id = f"msg-{i}"
payload = json.dumps({
"sensor_id": f"sensor-{random.randint(1, 5)}",
"temperature": random.uniform(20, 35),
"humidity": random.uniform(30, 80),
"timestamp": time.time()
})
success = process_message(message_id, payload)
print(f"Message {message_id} processed: {'success' if success else 'failure'}")
time.sleep(0.5) # Simulate the interval between messages
# Run the processor
process_stream_with_tracing()
The field of real-time analysis with Python keeps evolving rapidly. Some notable trends for 2025 include:
Integrating AI models into streaming pipelines is transforming real-time analysis.
Simply monitoring data is no longer enough; the demand now is for intelligent interpretation of events. Machine Learning and Deep Learning models are being built directly into streaming pipelines, using libraries such as torch, transformers, and onnxruntime for real-time inference.
# Conceptual example of object detection on a video stream
import cv2
import numpy as np
from confluent_kafka import Consumer, Producer
import json
import base64
import time
from transformers import DetrImageProcessor, DetrForObjectDetection
import torch
from PIL import Image
import io
# Load the object detection model
processor = DetrImageProcessor.from_pretrained("facebook/detr-resnet-50")
model = DetrForObjectDetection.from_pretrained("facebook/detr-resnet-50")
# Configure the Kafka consumer
consumer = Consumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'video-analyzer',
'auto.offset.reset': 'latest'
})
consumer.subscribe(['video-frames'])
# Configure the Kafka producer
producer = Producer({
'bootstrap.servers': 'localhost:9092'
})
# Function that detects objects
def detect_objects(image_bytes):
try:
# Convert bytes to an image
image = Image.open(io.BytesIO(image_bytes))
# Prepare the image for the model
inputs = processor(images=image, return_tensors="pt")
# Run the prediction
with torch.no_grad():
outputs = model(**inputs)
# Process the results
target_sizes = torch.tensor([image.size[::-1]])
results = processor.post_process_object_detection(
outputs, target_sizes=target_sizes, threshold=0.7
)[0]
detections = []
for score, label, box in zip(results["scores"], results["labels"], results["boxes"]):
box = [round(i, 2) for i in box.tolist()]
detections.append({
"label": model.config.id2label[label.item()],
"score": round(score.item(), 3),
"box": box
})
return detections
except Exception as e:
print(f"Object detection error: {e}")
return []
# Process the video stream
try:
while True:
msg = consumer.poll(1.0)
if msg is None:
continue
if msg.error():
print(f"Error: {msg.error()}")
continue
try:
# Decode the message
frame_data = json.loads(msg.value())
# Extract metadata and image
frame_id = frame_data.get('frame_id')
timestamp = frame_data.get('timestamp')
camera_id = frame_data.get('camera_id')
image_b64 = frame_data.get('image')
if not image_b64:
continue
# Decode the image
image_bytes = base64.b64decode(image_b64)
# Detect objects
start_time = time.time()
detections = detect_objects(image_bytes)
processing_time = time.time() - start_time
# Build the result
result = {
'frame_id': frame_id,
'timestamp': timestamp,
'camera_id': camera_id,
'detections': detections,
'processing_time': round(processing_time, 3),
'processed_at': time.time()
}
# Send the result to another topic
producer.produce(
'video-analysis',
key=str(frame_id),
value=json.dumps(result)
)
# Print the result
objects_found = [f"{d['label']} ({d['score']:.2f})" for d in detections]
print(f"Frame {frame_id}, Camera {camera_id}: {len(detections)} objects detected: {', '.join(objects_found)}")
# Flush the buffer periodically
producer.poll(0)
except json.JSONDecodeError:
print("Error decoding JSON")
except KeyboardInterrupt:
pass
finally:
consumer.close()
producer.flush()
Processing at the edge is becoming more important for reducing latency.
With the growing number of connected devices, moving part of the processing closer to the data source (such as IoT gateways or smart sensors) reduces latency and bandwidth.
# Conceptual example of edge processing with cloud synchronization
import numpy as np
import pandas as pd
import time
import json
import requests
from datetime import datetime
import threading
import queue
# Simulate an IoT sensor at the edge
class EdgeDevice:
def __init__(self, device_id, sync_interval=60):
self.device_id = device_id
self.sync_interval = sync_interval
self.local_buffer = []
self.cloud_queue = queue.Queue()
self.last_sync = time.time()
self.running = True
# Start the threads
self.sensor_thread = threading.Thread(target=self._sensor_loop)
self.processing_thread = threading.Thread(target=self._processing_loop)
self.sync_thread = threading.Thread(target=self._sync_loop)
def start(self):
print(f"Device {self.device_id} starting...")
self.sensor_thread.start()
self.processing_thread.start()
self.sync_thread.start()
def stop(self):
print(f"Device {self.device_id} stopping...")
self.running = False
self.sensor_thread.join()
self.processing_thread.join()
self.sync_thread.join()
print(f"Device {self.device_id} stopped.")
def _sensor_loop(self):
"""Simulates sensor readings"""
while self.running:
# Simulate a sensor reading
reading = {
'device_id': self.device_id,
'timestamp': datetime.now().isoformat(),
'temperature': np.random.uniform(20, 35),
'humidity': np.random.uniform(30, 80),
'pressure': np.random.uniform(980, 1020)
}
# Add to the local buffer
self.local_buffer.append(reading)
# Simulate the interval between readings
time.sleep(1)
def _processing_loop(self):
"""Processes data locally"""
while self.running:
if len(self.local_buffer) >= 10:
# Copy the buffer for processing
to_process = self.local_buffer.copy()
# Process locally
df = pd.DataFrame(to_process)
# Detect anomalies locally (simplified example)
mean_temp = df['temperature'].mean()
std_temp = df['temperature'].std()
for reading in to_process:
# Flag anomalies
temp = reading['temperature']
if abs(temp - mean_temp) > 2 * std_temp:
reading['anomaly'] = True
print(f"Anomaly detected locally! Temperature: {temp:.2f}°C")
# Send anomalies to the cloud immediately
self.cloud_queue.put(reading)
else:
reading['anomaly'] = False
# Calculate statistics
stats = {
'device_id': self.device_id,
'timestamp': datetime.now().isoformat(),
'window_start': to_process[0]['timestamp'],
'window_end': to_process[-1]['timestamp'],
'reading_count': len(to_process),
'temperature_mean': float(mean_temp),
'temperature_min': float(df['temperature'].min()),
'temperature_max': float(df['temperature'].max()),
'humidity_mean': float(df['humidity'].mean()),
'pressure_mean': float(df['pressure'].mean()),
'anomalies_detected': int(df['anomaly'].sum() if 'anomaly' in df else 0)
}
# Add the statistics to the sync queue
self.cloud_queue.put(stats)
time.sleep(1)
def _sync_loop(self):
"""Synchronizes data with the cloud"""
while self.running:
current_time = time.time()
# Sync periodically or when there is a lot of data
if current_time - self.last_sync >= self.sync_interval or self.cloud_queue.qsize() > 50:
batch = []
# Collect data from the queue
try:
while not self.cloud_queue.empty() and len(batch) < 100:
batch.append(self.cloud_queue.get_nowait())
except queue.Empty:
pass
if batch:
# Send to the cloud
try:
print(f"Syncing {len(batch)} items with the cloud...")
# In a real case, this would be an API call
# response = requests.post(
# "https://api.example.com/device-data",
# json={'device_id': self.device_id, 'data': batch},
# timeout=10
# )
# response.raise_for_status()
# Simulate a successful upload
time.sleep(0.5)
print(f"Sync complete: {len(batch)} items sent")
self.last_sync = current_time
except Exception as e:
print(f"Sync error: {e}")
# Put the items back in the queue
for item in batch:
self.cloud_queue.put(item)
time.sleep(1)
# Create and start the edge device
edge_device = EdgeDevice("sensor-edge-01")
try:
edge_device.start()
# Run for a while
time.sleep(120)
finally:
edge_device.stop()
Declarative languages are simplifying real-time analysis.
Frameworks such as Apache Flink, Beam, and KSQLDB are popularizing the use of continuous SQL queries to analyze data in real time. The trend is for less code to be needed to build complex pipelines, which increases productivity.
# Example with SQLStreamBuilder (conceptual)
from sqlstreambuilder import StreamBuilder, Schema, Field, WindowType
import time
# Define the schema
sensor_schema = Schema([
Field("sensor_id", "STRING"),
Field("temperature", "DOUBLE"),
Field("humidity", "DOUBLE"),
Field("timestamp", "TIMESTAMP")
])
# Create the builder
builder = StreamBuilder()
# Register the data source
builder.register_stream(
name="sensor_data",
schema=sensor_schema,
source={
"type": "kafka",
"topic": "sensor-data",
"bootstrap.servers": "localhost:9092",
"group.id": "sql-processor"
}
)
# Define the SQL queries
builder.create_view(
name="temperature_stats",
sql="""
SELECT
sensor_id,
TUMBLE_START(timestamp, INTERVAL '1' MINUTE) AS window_start,
TUMBLE_END(timestamp, INTERVAL '1' MINUTE) AS window_end,
COUNT(*) AS reading_count,
AVG(temperature) AS avg_temp,
MAX(temperature) AS max_temp,
MIN(temperature) AS min_temp
FROM sensor_data
GROUP BY
TUMBLE(timestamp, INTERVAL '1' MINUTE),
sensor_id
"""
)
builder.create_view(
name="humidity_stats",
sql="""
SELECT
sensor_id,
TUMBLE_START(timestamp, INTERVAL '1' MINUTE) AS window_start,
TUMBLE_END(timestamp, INTERVAL '1' MINUTE) AS window_end,
AVG(humidity) AS avg_humidity
FROM sensor_data
GROUP BY
TUMBLE(timestamp, INTERVAL '1' MINUTE),
sensor_id
"""
)
builder.create_view(
name="high_temperature_alerts",
sql="""
SELECT
sensor_id,
temperature,
timestamp
FROM sensor_data
WHERE temperature > 30.0
"""
)
# Define the outputs
builder.create_sink(
name="stats_sink",
source="temperature_stats",
sink={
"type": "kafka",
"topic": "temperature-stats",
"bootstrap.servers": "localhost:9092"
}
)
builder.create_sink(
name="humidity_sink",
source="humidity_stats",
sink={
"type": "kafka",
"topic": "humidity-stats",
"bootstrap.servers": "localhost:9092"
}
)
builder.create_sink(
name="alerts_sink",
source="high_temperature_alerts",
sink={
"type": "kafka",
"topic": "temperature-alerts",
"bootstrap.servers": "localhost:9092"
}
)
# Run the pipeline
job = builder.build()
job.start()
try:
while True:
time.sleep(1)
except KeyboardInterrupt:
job.stop()
The convergence of batch and real-time processing is simplifying architectures.
Batch and stream are being unified into a single processing architecture. Tools such as Apache Beam and Delta Lake are making this hybrid model easier: we write the code once and run it in different execution modes.
# Conceptual example with Apache Beam for unified processing
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.transforms.window import FixedWindows
from apache_beam.io.kafka import ReadFromKafka, WriteToKafka
import json
import datetime
# Define the transformations
class ParseJson(beam.DoFn):
def process(self, element):
try:
record = json.loads(element.decode('utf-8'))
return [record]
except Exception as e:
print(f"Error parsing JSON: {e}")
return []
class AddTimestamp(beam.DoFn):
def process(self, element):
# Add a timestamp for windowing
timestamp = element.get('timestamp')
if timestamp:
# Convert to a Beam timestamp
event_time = datetime.datetime.fromtimestamp(timestamp)
yield beam.window.TimestampedValue(element, event_time.timestamp())
else:
yield element
class CalculateStats(beam.DoFn):
def process(self, window_pair):
key, readings = window_pair
if not readings:
return []
temperatures = [r.get('temperature', 0) for r in readings if 'temperature' in r]
humidities = [r.get('humidity', 0) for r in readings if 'humidity' in r]
if not temperatures:
return []
stats = {
'sensor_id': key,
'window_timestamp': datetime.datetime.now().isoformat(),
'reading_count': len(readings),
'temperature_avg': sum(temperatures) / len(temperatures),
'temperature_min': min(temperatures),
'temperature_max': max(temperatures)
}
if humidities:
stats['humidity_avg'] = sum(humidities) / len(humidities)
return [json.dumps(stats).encode('utf-8')]
# Configure the pipeline
pipeline_options = PipelineOptions([
'--runner=DirectRunner',
'--streaming'
])
with beam.Pipeline(options=pipeline_options) as pipeline:
# Read from Kafka
readings = (
pipeline
| 'ReadFromKafka' >> ReadFromKafka(
consumer_config={
'bootstrap.servers': 'localhost:9092',
'auto.offset.reset': 'latest'
},
topics=['sensor-data']
)
| 'ParseJson' >> beam.ParDo(ParseJson())
| 'AddEventTimestamps' >> beam.ParDo(AddTimestamp())
)
# Process in time windows
windowed_stats = (
readings
| 'WindowByMinute' >> beam.WindowInto(FixedWindows(60)) # 1-minute windows
| 'ExtractSensorId' >> beam.Map(lambda x: (x.get('sensor_id', 'unknown'), x))
| 'GroupBySensor' >> beam.GroupByKey()
| 'CalculateStats' >> beam.ParDo(CalculateStats())
)
# Detect anomalies
anomalies = (
readings
| 'FilterHighTemperature' >> beam.Filter(lambda x: x.get('temperature', 0) > 30.0)
| 'FormatAnomaly' >> beam.Map(lambda x: json.dumps({
'sensor_id': x.get('sensor_id', 'unknown'),
'temperature': x.get('temperature', 0),
'timestamp': x.get('timestamp', 0),
'alert_type': 'HIGH_TEMPERATURE',
'detected_at': datetime.datetime.now().isoformat()
}).encode('utf-8'))
)
# Write the results to Kafka
windowed_stats | 'WriteStatsToKafka' >> WriteToKafka(
producer_config={'bootstrap.servers': 'localhost:9092'},
topic='sensor-stats'
)
anomalies | 'WriteAnomaliesToKafka' >> WriteToKafka(
producer_config={'bootstrap.servers': 'localhost:9092'},
topic='sensor-anomalies'
)
Real-time data analysis with Python will keep evolving rapidly in the coming years. The tools and techniques presented in this article provide a solid foundation for implementing data streaming solutions that can process information in real time, extract valuable insights, and trigger immediate actions.
In 2025, organizations that master real-time analysis have a significant competitive advantage: they can respond quickly to market changes, detect anomalies before they become problems, and offer customers personalized experiences.
By combining the power and flexibility of Python with modern streaming technologies such as Kafka, Spark, and Flink, you can build robust real-time processing systems that scale to meet your business’s demands, regardless of the volume, velocity, or variety of your data.
Artificial intelligence is revolutionizing every sector of society, from enterprise applications to everyday solutions. At…
In this lesson of the Python mini course, I want to cover two collection types…
Among Python's most versatile data structures, dictionaries stand out. They help us represent structured data…
In this ninth lesson of the mini course, I want to talk about one of…
Understanding variable scope in Python is essential to avoid errors and write more predictable programs.…
Functions are one of the most important concepts in any programming language, and Python is…
Este blog utiliza cookies. Se você continuar assumiremos que você está satisfeito com ele.
Leia Mais...