【问题标题】:Catch only the payload of CDC in Pyspark structured streaming?在 Pyspark 结构化流中仅捕获 CDC 的有效负载?
【发布时间】:2021-09-13 00:54:35
【问题描述】:
  • 我正在尝试创建一条从 SQL Server 到 Pyspark 的管道以捕获 SQL Server 中的数据更改,我已准备好一切:
    • 在 SQL Server 中启用 CDC
    • 从 SQL Server 生产到 Kafka 并从 Pyspark 结构化流中的 Kafka 主题消费。
  • 问题是:当我尝试使用控制台消费者检查数据更改是否通过 Kafka 时,它会显示 JSON 格式的消息,分为两条记录:Schema 和 Payload,在 Payload 内部有 Before 和 After 给出你分别是更改前的数据和更改后的数据。
    • 我只是在负载中被处理-->在此 JSON 消息的一部分之后
      • 因为当我像这样流式传输它时,在 Jupyter 命令行中我需要的字段显示 null ,我理解这是因为 JSON 格式很复杂
    • 这是我的 pyspark 代码:
     import os

os.environ['PYSPARK_SUBMIT_ARGS'] = f'--packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2 pyspark-shell'

import findspark

findspark.init()

import pyspark
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark.sql.types import *
import time

kafka_topic_name = "test-spark"
kafka_bootstrap_servers = '192.168.1.3:9092'

spark = SparkSession \
    .builder \
    .appName("PySpark Structured Streaming with Kafka and Message Format as JSON") \
    .master("local[*]") \
    .getOrCreate()

# Construct a streaming DataFrame that reads from TEST-SPARK
df = spark \
    .readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", kafka_bootstrap_servers) \
    .option("subscribe", kafka_topic_name) \
    .load()

print("Printing Schema of df: ")
df.printSchema()


df1 = df.selectExpr("CAST(value AS STRING)", "timestamp")
df1.printSchema()

 schema = StructType() \
        .add("name", StringType()) \
        .add("type", StringType())

df2 = df1\
        .select(from_json(col("value"), schema)\
        .alias("records"), "timestamp")
    df3 = df2.select("records.*", "timestamp")

  print("Printing Schema of records_df3: ")
    df3.printSchema()

 records_write_stream = df3 \
        .writeStream \
        .trigger(processingTime='5 seconds') \
        .outputMode("update") \
        .option("truncate", "false")\
        .format("console") \
        .start()
    records_write_stream.awaitTermination()

    print("Stream Data Processing Application Completed.")
  • 这是一张显示 CDC 消息到达 Kafka 的图像:

  • 如果有人知道如何仅使用 Payload-->在参与 Pyspark 结构化流式传输之后,请帮助我。

【问题讨论】:

  • 请分享您的 pyspark 代码,从您连接到 kafka 流式实例时到您在问题中分享的位置
  • @ggordon 已经完成了,如果您有任何建议帮助,请。

标签: sql-server apache-spark pyspark apache-kafka cdc


【解决方案1】:

-在搜索了更多之后,我发现了如何仅显示和捕获 CDC msg 的有效负载部分。

  • 您需要将此添加到您的 Worker.properties:
value.converter=org.apache.kafka.connect.json.JsonConverter

value.converter.schemas.enable=false

【讨论】:

    【解决方案2】:

    您应该将您的 Debezeium 连接器修改为具有 value.converter.schemas.enabled=false,然后您将只需要使用 payload 字段。

    否则,您可以为整个对象创建一个类/模式以及 from_json() 函数,或者将值保留为字符串并使用 get_json_object() Spark 函数解析数据

    同样相关 - 您可能想要提取 NewRecordState

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-02-16
      • 2013-08-04
      • 1970-01-01
      • 2020-06-30
      • 2013-10-19
      • 1970-01-01
      相关资源
      最近更新 更多