Pulsar Schema使用原理介绍

2024-03-09 19:04

本文主要是介绍Pulsar Schema使用原理介绍,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

一、引言

关于Pulsar Schema,咱们要想想以下几个问题

  • Pulsar 中的 Schema 是什么?
  • Pulsar Schema Registry的作用是什么?
  • 怎么使用?
  • 原理是什么?

二、Schema 是什么

Schema是定义结构化数据和二进制字节数组之间转换的逻辑,Pulsar的消息是以非结构化的二进制数组进行存储的,Schema只有在读写时才会被应用于数据上,因此生产者和消费者需要对Schema达成一致。Pulsar通过Schema Registry作为一个中央仓库存储Schema信息,它可以协调生产者和消费者保证相同的Schema,它可以存储多个版本的Schema,支持不同的兼容性配置以及根据兼容性的要求进行Schema的演进。

Pulsar将Schema存储在Bookie上,Schema的写入、读取都通过Broker和Bookie交互,这个逻辑跟消息的读写操作是一只的,因此不需要额外考虑Schema的可用性和可靠性问题,因此整体看Pulsar实现Schema Registry的方式非常优雅

类型安全在所有数据应用中都非常重要,生产者和消费者需要某种机制协调数据类型来避免各种潜在的问题,比如序列化和反序列化方式不一致。数据安全通常有两种处理方式client-side和service-side,本质上就是客户端用时决定和服务端提前保证

client-side:将一切交给用户,客户端自行负责消息的序列化和反序列化并且保证生产消费时消息的类型安全,这种方式的最大问题就是类型是通过约定的,一旦生产者写入非约定的数据,下游的消费者将没有办法解析数据

server-side:数据安全由服务端保证,生产者和消费者都需要跟服务端提前确定数据类型。这种方式真正意义上保证了数据的类型安全,避免了生产者写入非法数据的问题

两种差异如下图
在这里插入图片描述

三、怎么使用

1. client-side

生产者代码逻辑

    //schema在第一次写入的时候就已经决定好了,后续用其他的schema消息类型会写入失败public static void customSchemaProducer()  {try {String serverUrl = "http://localhost:8080";PulsarClient pulsarClient =PulsarClient.builder().serviceUrl(serverUrl).build();Producer<byte[]> producer = pulsarClient.newProducer().topic("persistent://sherlock-api-tenant-1/sherlock-namespace-1/topic_8").create();User user = new User();user.setName("老六");user.setAge(21);user.setAddress("海南");//由用户自行做序列化逻辑producer.send(JSON.toJSONString(user).getBytes());producer.close();pulsarClient.close();} catch (Exception e) {}}

消费者逻辑代码

  public static void customSchemaConsumer() {try {String serverUrl = "http://localhost:8080";PulsarClient pulsarClient =PulsarClient.builder().serviceUrl(serverUrl).build();Consumer<byte[]> consumer = pulsarClient.newConsumer().topic("persistent://sherlock-api-tenant-1/sherlock-namespace-1/topic_8").subscriptionName("sub_03").subscribe();while(true) {Message<byte[]> message = consumer.receive();//由用户自行做序列化逻辑byte[] user = message.getValue();System.out.println("消息数据为:"+JSON.parseObject(user, User.class).toString());consumer.acknowledge(message);//consumer.negativeAcknowledge(message);}} catch (Exception e) {}}

执行效果如下

消息数据为:User{name='老六', age=21, address='海南'}

2. server-side

生产者代码逻辑

   //schema在第一次写入的时候就已经决定好了,后续用其他的schema消息类型会写入失败public static void customSchemaProducer()  {try {String serverUrl = "http://localhost:8080";PulsarClient pulsarClient =PulsarClient.builder().serviceUrl(serverUrl).build();//由Pulsar做序列化逻辑Producer<User> producer = pulsarClient.newProducer(AvroSchema.of(User.class)).topic("persistent://sherlock-api-tenant-1/sherlock-namespace-1/topic_7").create();User user = new User();user.setName("王武");user.setAge(36);user.setAddress("海南");producer.send(user);producer.close();pulsarClient.close();} catch (Exception e) {}}

消费者逻辑代码

   public static void customSchemaConsumer() {try {String serverUrl = "http://localhost:8080";PulsarClient pulsarClient =PulsarClient.builder().serviceUrl(serverUrl).build();//由Pulsar做序列化逻辑Consumer<User> consumer = pulsarClient.newConsumer(AvroSchema.of(User.class)).topic("persistent://sherlock-api-tenant-1/sherlock-namespace-1/topic_7").subscriptionName("sub_03").subscribe();while(true) {Message<User> message = consumer.receive();User user = message.getValue();System.out.println("消息数据为:"+user);consumer.acknowledge(message);//consumer.negativeAcknowledge(message);}} catch (Exception e) {}}

执行效果如下

消息数据为:User{name='王武', age=36, address='海南'}

3. 小结

分别查看两个Topic的Schema信息如下图

通过查询client-side的Schema信息,会发现Pulsar服务端其实并没有进行存储,相当于不指定Schema的话Pulsar默认都用byte数组
在这里插入图片描述

再来看看server-side的Schema信息,可以看到打印如下,namespace是pojo类的包路径,name是pojo类名,然后fields就是pojo类的各个字段的属性(像不像mysql里面的表结构,不少场景Topic就是当作表来用的),然后type是AVRO是由于咱们是用的avro进行序列化的。
在这里插入图片描述

除了在读写数据时指定Schema,Pulsar还支持通过admin管理流提前指定好,具体指令在这里。如果是用Pulsar来作为实时数仓场景,强烈建议提前通过admin管理流进行指定好,配置isSchemaValidationEnforced可以考虑开启。如果条件允许可以考虑做成服务化,例如通过Web页面提供新建Schema、修改Schema操作并接入公司内部的审批流等

pulsar-admin schemas upload --filename
POST /admin/v2/schemas/:tenant/:namespace/:topic/schema
pulsar-admin schemas get sherlock-api-tenant-1/sherlock-namespace-1/partition_partition_topic_3

四、原理解析

Schema相关的流程咱们需要关注以下几个

  • 注册Schema流程
    • 生产者端侧
    • 消费者端侧
    • 指定服务器
  • Schema生效流程
  • 更新Schema流程

1. 注册Schema流程

生产者端侧
在这里插入图片描述

  1. 生产者实例会在内部构造schema实例,生产者会通过它对数据进行转换
  2. 生产者会请求连接Broker,并传递schema信息 SchemaInfo
  3. Broker会在schema registry中查找这个schema是否被注册,如果已经注册了就将注册的schema版本返回给生产者
  4. Broker检查是否支持自动更新schema,如果配置不允许自动更新,则这个schema不能被注册并且拒绝生产者
  5. Broker进行schema兼容性检查,如果通过检查则将此schema存储在schema registry并返回版本给生产者,生产者所有消息以这个schema格式进行发送;若是检查没通过则拒绝生产者

消费者端
在这里插入图片描述

  1. 消费者实例会在内部构造schema实例
  2. 消费者请求连接Broker,并传递schema信息 SchemaInfo
  3. Broker检查这个Topic是否已经在使用,有的话跳到第五步,否则跳到第四步
  4. Broker检查是否支持自动更新schema,如果支持则注册这个schema,否则拒绝客户端
  5. Broker进行schema兼容性检查,通过则连接否则拒绝客户端

五、参考文献

  1. Pulsar:Schema Registry介绍
  2. 官方文档
  3. 深度解读 Pulsar Schema

这篇关于Pulsar Schema使用原理介绍的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/791658

相关文章

postgresql使用UUID函数的方法

《postgresql使用UUID函数的方法》本文给大家介绍postgresql使用UUID函数的方法,本文给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的朋友参考下吧... 目录PostgreSQL有两种生成uuid的方法。可以先通过sql查看是否已安装扩展函数,和可以安装的扩展函数

如何使用Lombok进行spring 注入

《如何使用Lombok进行spring注入》本文介绍如何用Lombok简化Spring注入,推荐优先使用setter注入,通过注解自动生成getter/setter及构造器,减少冗余代码,提升开发效... Lombok为了开发环境简化代码,好处不用多说。spring 注入方式为2种,构造器注入和setter

MySQL中比较运算符的具体使用

《MySQL中比较运算符的具体使用》本文介绍了SQL中常用的符号类型和非符号类型运算符,符号类型运算符包括等于(=)、安全等于(=)、不等于(/!=)、大小比较(,=,,=)等,感兴趣的可以了解一下... 目录符号类型运算符1. 等于运算符=2. 安全等于运算符<=>3. 不等于运算符<>或!=4. 小于运

使用zip4j实现Java中的ZIP文件加密压缩的操作方法

《使用zip4j实现Java中的ZIP文件加密压缩的操作方法》本文介绍如何通过Maven集成zip4j1.3.2库创建带密码保护的ZIP文件,涵盖依赖配置、代码示例及加密原理,确保数据安全性,感兴趣的... 目录1. zip4j库介绍和版本1.1 zip4j库概述1.2 zip4j的版本演变1.3 zip4

Python 字典 (Dictionary)使用详解

《Python字典(Dictionary)使用详解》字典是python中最重要,最常用的数据结构之一,它提供了高效的键值对存储和查找能力,:本文主要介绍Python字典(Dictionary)... 目录字典1.基本特性2.创建字典3.访问元素4.修改字典5.删除元素6.字典遍历7.字典的高级特性默认字典

使用Python构建一个高效的日志处理系统

《使用Python构建一个高效的日志处理系统》这篇文章主要为大家详细讲解了如何使用Python开发一个专业的日志分析工具,能够自动化处理、分析和可视化各类日志文件,大幅提升运维效率,需要的可以了解下... 目录环境准备工具功能概述完整代码实现代码深度解析1. 类设计与初始化2. 日志解析核心逻辑3. 文件处

一文详解如何使用Java获取PDF页面信息

《一文详解如何使用Java获取PDF页面信息》了解PDF页面属性是我们在处理文档、内容提取、打印设置或页面重组等任务时不可或缺的一环,下面我们就来看看如何使用Java语言获取这些信息吧... 目录引言一、安装和引入PDF处理库引入依赖二、获取 PDF 页数三、获取页面尺寸(宽高)四、获取页面旋转角度五、判断

C++中assign函数的使用

《C++中assign函数的使用》在C++标准模板库中,std::list等容器都提供了assign成员函数,它比操作符更灵活,支持多种初始化方式,下面就来介绍一下assign的用法,具有一定的参考价... 目录​1.assign的基本功能​​语法​2. 具体用法示例​​​(1) 填充n个相同值​​(2)

Spring StateMachine实现状态机使用示例详解

《SpringStateMachine实现状态机使用示例详解》本文介绍SpringStateMachine实现状态机的步骤,包括依赖导入、枚举定义、状态转移规则配置、上下文管理及服务调用示例,重点解... 目录什么是状态机使用示例什么是状态机状态机是计算机科学中的​​核心建模工具​​,用于描述对象在其生命

使用Python删除Excel中的行列和单元格示例详解

《使用Python删除Excel中的行列和单元格示例详解》在处理Excel数据时,删除不需要的行、列或单元格是一项常见且必要的操作,本文将使用Python脚本实现对Excel表格的高效自动化处理,感兴... 目录开发环境准备使用 python 删除 Excphpel 表格中的行删除特定行删除空白行删除含指定