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

无需Spark:从Kafka流数据提取部分内容存入Pandas DataFrame

Storing Kafka Stream Data into a Pandas DataFrame

Got it, let's walk through modifying your Confluent Kafka Consumer code to capture parts of your stream data into a Pandas DataFrame. Here's a practical, step-by-step solution:

Step 1: Add Required Imports

First, you'll need to import pandas (for DataFrame handling) and json (since Kafka messages are commonly serialized as JSON—adjust this if you're using a different format like Avro).

Step 2: Use an Efficient Data Collection Strategy

Instead of appending directly to a DataFrame (which is slow for large streams), we'll use a list to collect our target data points first. We can convert this list to a DataFrame later—either after hitting a message limit, on a schedule, or when stopping the consumer.

Modified Working Code

from confluent_kafka import Consumer, KafkaError
import pandas as pd
import json

# Initialize Kafka Consumer with your configs
c = Consumer({
    'bootstrap.servers': "###",
    'group.id': '###',
    'default.topic.config': {
        'auto.offset.reset': 'latest'
    }
})
c.subscribe(['scorestore'])

# Empty list to collect our desired stream data points
data_list = []

# Optional: Set a message limit to avoid infinite streaming (remove if you want continuous processing)
message_limit = 1000
message_count = 0

try:
    while True:
        msg = c.poll(1.0)
        if msg is None:
            continue
        if msg.error():
            if msg.error().code() == KafkaError._PARTITION_EOF:
                continue
            else:
                print(f"Error encountered: {msg.error()}")
                break
        
        # Parse the Kafka message (adjust this based on your actual serialization format)
        # Assuming message value is UTF-8 encoded JSON
        message_content = json.loads(msg.value().decode('utf-8'))
        
        # Extract the specific fields you want to store (replace with your actual fields)
        # Example: grabbing 'user_id', 'score', and 'event_time' from the message
        filtered_data = {
            'user_id': message_content.get('user_id'),
            'score': message_content.get('score'),
            'event_time': message_content.get('event_time')
        }
        
        # Add to our collection list
        data_list.append(filtered_data)
        message_count += 1
        
        # Print progress (optional)
        print(f"Processed message {message_count}: {filtered_data}")
        
        # Stop processing once we hit our limit (remove this block for non-stop streaming)
        if message_count >= message_limit:
            print(f"Reached message limit of {message_limit}, stopping consumer.")
            break

finally:
    # Always clean up the consumer connection
    c.close()
    
    # Convert collected data to a Pandas DataFrame
    stream_df = pd.DataFrame(data_list)
    print("\nFinal Data Preview:")
    print(stream_df.head())
    
    # Optional: Save the DataFrame to a file (CSV, Parquet, etc.)
    # stream_df.to_csv('kafka_stream_results.csv', index=False)

Key Tips & Notes:

  • Message Parsing: If your messages use a different serialization format (like Avro or Protobuf), swap out the json.loads step with the appropriate deserializer (e.g., confluent_kafka.schema_registry.avro.AvroDeserializer for Avro).
  • Efficiency: Using a list to batch-collect data is way more performant than appending rows one-by-one to a DataFrame, especially for large streams.
  • Continuous Streaming: For non-stop processing, remove the message limit block and add logic to periodically convert the list to a DataFrame (e.g., every 500 messages) and clear the list to avoid memory bloat.
  • Error Safety: The try/finally block ensures the consumer is properly closed even if the script is interrupted.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:30:36