如何使用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:
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).
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()
- 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 aCacheEntryEventwith 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.
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.
- 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

