Introduction #
SmartChart supports Kafka as a data source for real-time scenarios, retrieving the latest message from a specified partition. Uses custom data source functions to consume the latest record from a Kafka Topic.
| Key Concept | Description |
|---|---|
dataset() |
Returns query dataset, reads latest Kafka record |
insert_dataset() |
Data submission implementation, writes messages to Kafka |
| Consumer config | Uses KafkaConsumer with SASL authentication |
| Partition | Configured via config['db'] |
Usage #
Refer to “Custom Data Source” for usage. Below is reference code:
def dataset(*args, **kwargs):
"""
Returns query dataset
:return: 2D array or JSON dict
"""
from kafka import KafkaConsumer, TopicPartition
import json
sqlList = args[0] # Dataset editor input split by semicolons [sql1, sql2...]
config = args[1] # Config dict {'host','port','user','password','db'}
result = {}
consumer = KafkaConsumer(sasl_mechanism='PLAIN',
security_protocol='SASL_PLAINTEXT',
sasl_plain_username=config['user'],
sasl_plain_password=config['password'],
bootstrap_servers=config['host'],
auto_offset_reset='earliest',
api_version=(1, 0, 0),
consumer_timeout_ms=50,
value_deserializer=lambda v: json.loads(v.decode('utf-8')),
)
topic = sqlList[0]
partition = int(config['db'])
tp = TopicPartition(topic=topic, partition=partition)
consumer.assign([tp])
end_offsets = consumer.end_offsets([tp]).get(tp)
consumer.seek(tp, offset=end_offsets-1)
for message in consumer:
result = message.value
break
return result
def insert_dataset(*args, **kwargs):
"""
Data submission implementation
"""
from kafka import KafkaProducer
import json
contents = args[0]
table = args[1]
config = args[3]
producer = KafkaProducer(sasl_mechanism='PLAIN',
security_protocol='SASL_PLAINTEXT',
sasl_plain_username=config['user'],
sasl_plain_password=config['password'],
bootstrap_servers=config['host'],
value_serializer=lambda v: json.dumps(v).encode('utf-8')
)
producer.send(table, value=contents, partition=int(config['db']))