Kafka Source

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']))