Options
- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
08-31-2023 04:12 PM
We've now added a native connector with parsing directly with Spark Dataframes.
https://docs.databricks.com/en/structured-streaming/protocol-buffers.html
from pyspark.sql.protobuf.functions import to_protobuf, from_protobuf schema_registry_options = { "schema.registry.subject" : "app-events-value", "schema.registry.address" : "https://schema-registry:8081/" } # Convert binary Protobuf to SQL struct with from_protobuf(): proto_events_df = ( input_df .select( from_protobuf("proto_bytes", options = schema_registry_options) .alias("proto_event") ) ) # Convert SQL struct to binary Protobuf with to_protobuf(): protobuf_binary_df = ( proto_events_df .selectExpr("struct(name, id, context) as event") .select( to_protobuf("event", options = schema_registry_options) .alias("proto_bytes") ) )