| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
| Name | Name | Last commit date | ||
|---|---|---|---|---|
parent directory.. | ||||
Python SDK for Streamdal.
For more details, see the main streamdal repo.
See https://docs.streamdal.com
python -m pip install streamdal
import json
from streamdal import (OPERATION_TYPE_CONSUMER, ProcessRequest, StreamdalClient, StreamdalConfig, EXEC_STATUS_TRUE)
client = StreamdalClient(
cfg=StreamdalConfig(
service_name="order-ingest",
streamdal_url="streamdal-server.svc.cluster.local:8082",
streamdal_token="1234",
)
)
res = client.process(
ProcessRequest(
operation_type=OPERATION_TYPE_CONSUMER,
operation_name="new-order-topic",
component_name="kafka",
data=b'{"object": {"email": "user@streamdal.com"}}',
)
)
# Check that process() completed successfully
if res.status == EXEC_STATUS_TRUE:
print("Success processed payload")
data = json.loads(res.data)
print("Response:", json.dumps(data, indent=2))
else:
print("Failed to process payload")
print("Error:", res.status_message)Metrics are published to Streamdal server and are available in Prometheus format at http://streamdal_server_url:8081/metrics
| Metric | Description | Labels |
|---|---|---|
| streamdal_counter_consume_bytes | Number of bytes consumed by the client | service, component_name, operation_name, pipeline_id, pipeline_name |
| streamdal_counter_consume_errors | Number of errors encountered while consuming payloads | service, component_name, operation_name, pipeline_id, pipeline_name |
| streamdal_counter_consume_processed | Number of payloads processed by the client | service, component_name, operation_name, pipeline_id, pipeline_name |
| streamdal_counter_produce_bytes | Number of bytes produced by the client | service, component_name, operation_name, pipeline_id, pipeline_name |
| streamdal_counter_produce_errors | Number of errors encountered while producing payloads | service, component_name, operation_name, pipeline_id, pipeline_name |
| streamdal_counter_produce_processed | Number of payloads processed by the client | service, component_name, operation_name, pipeline_id, pipeline_name |
| streamdal_counter_notify | Number of notifications sent to the server | service, component_name, operation_name, pipeline_id, pipeline_name |
Any push or merge to the main branch with any changes in /sdks/python/* will automatically tag and release a new console version with sdks/python/vX.Y.Z.
(1) If you'd like to skip running the release action on push/merge to main, include norelease anywhere in the commit message.
| Back | FazBrowse Home | New Git URL |