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

如何使用Python实现Apache Ignite连续查询?已查阅Java相关示例

Great question! Apache Ignite's Python client does support continuous queries—even though the official docs lean heavily on Java examples, the Python API mirrors much of that functionality. Here's a step-by-step guide with a complete working example:

Prerequisites

First, ensure you have the Ignite Python client installed:

pip install apache-ignite

Also, confirm your Ignite cluster is running and accessible (default port is 10800).

Basic Continuous Query Example

This example sets up a continuous query that listens for cache changes and prints event details whenever an entry is created, updated, or deleted:

from ignite import IgniteClient
from ignite.datatypes import CacheEntryEvent
from ignite.events import EventType

def handle_event(event: CacheEntryEvent):
    """Callback function to process cache events"""
    event_label = {
        EventType.EVT_CACHE_ENTRY_CREATED: "CREATED",
        EventType.EVT_CACHE_ENTRY_UPDATED: "UPDATED",
        EventType.EVT_CACHE_ENTRY_DELETED: "DELETED"
    }.get(event.event_type, "UNKNOWN")
    
    print(f"Event: {event_label} | Key: {event.key} | Old Value: {event.old_value} | New Value: {event.new_value}")

if __name__ == "__main__":
    # Connect to the Ignite cluster
    client = IgniteClient()
    client.connect("127.0.0.1", 10800)
    
    # Get or create a cache instance
    cache = client.get_or_create_cache("my_continuous_cache")
    
    # Initialize the continuous query
    query = client.create_continuous_query()
    
    # Attach the event listener (triggers when events are received)
    query.set_listener(handle_event)
    
    # Optional: Filter events server-side to reduce network traffic
    # query.set_remote_filter(lambda evt: evt.event_type in [EventType.EVT_CACHE_ENTRY_CREATED, EventType.EVT_CACHE_ENTRY_UPDATED])
    
    # Start listening for events
    query.start(cache)
    
    try:
        # Perform cache operations to trigger events
        cache.put("user_1", {"name": "Alice", "age": 30})
        cache.put("user_1", {"name": "Alice", "age": 31})
        cache.remove("user_1")
        
        # Keep the program running to listen for ongoing events
        input("Press Enter to stop listening...\n")
    finally:
        # Cleanup: Stop the query and close the client connection
        query.stop()
        client.close()
Key Components Explained
  • IgniteClient: Establishes and manages the connection to your Ignite cluster.
  • ContinuousQuery: The core object defining your continuous query. Configure it with:
    • set_listener(): A callback function that processes incoming events. It receives a CacheEntryEvent with details like event type, key, old/new values.
    • set_remote_filter(): An optional server-side filter to only send relevant events (e.g., keys matching a pattern). This reduces unnecessary network traffic.
  • EventType: Enumeration of cache event types (CREATED, UPDATED, DELETED, etc.) to target specific changes.
Advanced: Server-Side Filtering

If you want to only receive events for keys matching a specific condition (e.g., keys starting with "user_"), use a remote filter:

def user_key_filter(event: CacheEntryEvent):
    # Only forward events for keys starting with "user_"
    return isinstance(event.key, str) and event.key.startswith("user_")

query.set_remote_filter(user_key_filter)

Note: Remote filters must be serializable, so avoid functions with external dependencies or non-serializable variables.

Important Notes
  • Always call query.stop() when finished to free cluster resources.
  • For custom data types, ensure they're properly serialized/deserialized by the Python client (you may need to register custom serializers).
  • The Python API aligns with the Java API, so you can reference official docs for advanced configs like event batching size or interval.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 08:43:51