1

Python 2.7 と Apache Avro(python client) を使用して、kafka ブローカーを介してシリアル化されたメッセージを交換しようとしています。以前にスキーマを作成せずにメッセージを交換する方法があるかどうか知りたいです。

これはコードです(スキーマ、sensor.avscを使用して、私が避けたいことです):

from kafka import SimpleProducer, KafkaClient
import avro.schema
import io, random
from avro.io import DatumWriter

# To send messages synchronously
kafka = KafkaClient('localhost:9092')
producer = SimpleProducer(kafka, async = False)

# Kafka topic
topic = "sensor_network_01"

# Path to user.avsc avro schema that i don't want
schema_path="sensor.avsc"
schema = avro.schema.parse(open(schema_path).read())


for i in xrange(100):
    writer = avro.io.DatumWriter(schema)
    bytes_writer = io.BytesIO()
    encoder = avro.io.BinaryEncoder(bytes_writer)
    # creation of random data
    writer.write({"sensor_network_name": "Sensor_1", "value": random.randint(0,10), "threshold_value":10 }, encoder)

    raw_bytes = bytes_writer.getvalue()
    producer.send_messages(topic, raw_bytes)

これは、sensor.avsc ファイルです。

{
    "namespace": "sensors.avro",
    "type": "record",
    "name": "Sensor",
    "fields": [
        {"name": "sensor_network_name", "type": "string"},
        {"name": "value",  "type": ["int", "null"]},
        {"name": "threshold_value", "type": ["int", "null"]}
    ]
}
4

2 に答える 2

3

このコード:

import avro.schema
import io, random
from avro.io import DatumWriter, DatumReader
import avro.io

# Path to user.avsc avro schema
schema_path="user.avsc"
schema = avro.schema.Parse(open(schema_path).read())


for i in xrange(1):
    writer = avro.io.DatumWriter(schema)
    bytes_writer = io.BytesIO()
    encoder = avro.io.BinaryEncoder(bytes_writer)
    writer.write({"name": "123", "favorite_color": "111", "favorite_number": random.randint(0,10)}, encoder)
    raw_bytes = bytes_writer.getvalue()

    print(raw_bytes)

    bytes_reader = io.BytesIO(raw_bytes)
    decoder = avro.io.BinaryDecoder(bytes_reader)
    reader = avro.io.DatumReader(schema)
    user1 = reader.read(decoder)
    print(" USER = {}".format(user1))

このスキーマを扱うため

{"namespace": "example.avro",
 "type": "record",
 "name": "User",
 "fields": [
     {"name": "name", "type": "string"},
     {"name": "favorite_number",  "type": ["int", "null"]},
     {"name": "favorite_color", "type": ["string", "null"]}
 ]
}

必要なものです。

クレジットはこの要点に行きます

于 2018-01-23T17:01:43.580 に答える