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

相关文章

python使用fastapi实现多语言国际化的操作指南

《python使用fastapi实现多语言国际化的操作指南》本文介绍了使用Python和FastAPI实现多语言国际化的操作指南,包括多语言架构技术栈、翻译管理、前端本地化、语言切换机制以及常见陷阱和... 目录多语言国际化实现指南项目多语言架构技术栈目录结构翻译工作流1. 翻译数据存储2. 翻译生成脚本

C++ Primer 多维数组的使用

《C++Primer多维数组的使用》本文主要介绍了多维数组在C++语言中的定义、初始化、下标引用以及使用范围for语句处理多维数组的方法,具有一定的参考价值,感兴趣的可以了解一下... 目录多维数组多维数组的初始化多维数组的下标引用使用范围for语句处理多维数组指针和多维数组多维数组严格来说,C++语言没

在 Spring Boot 中使用 @Autowired和 @Bean注解的示例详解

《在SpringBoot中使用@Autowired和@Bean注解的示例详解》本文通过一个示例演示了如何在SpringBoot中使用@Autowired和@Bean注解进行依赖注入和Bean... 目录在 Spring Boot 中使用 @Autowired 和 @Bean 注解示例背景1. 定义 Stud

使用 sql-research-assistant进行 SQL 数据库研究的实战指南(代码实现演示)

《使用sql-research-assistant进行SQL数据库研究的实战指南(代码实现演示)》本文介绍了sql-research-assistant工具,该工具基于LangChain框架,集... 目录技术背景介绍核心原理解析代码实现演示安装和配置项目集成LangSmith 配置(可选)启动服务应用场景

使用Python快速实现链接转word文档

《使用Python快速实现链接转word文档》这篇文章主要为大家详细介绍了如何使用Python快速实现链接转word文档功能,文中的示例代码讲解详细,感兴趣的小伙伴可以跟随小编一起学习一下... 演示代码展示from newspaper import Articlefrom docx import

oracle DBMS_SQL.PARSE的使用方法和示例

《oracleDBMS_SQL.PARSE的使用方法和示例》DBMS_SQL是Oracle数据库中的一个强大包,用于动态构建和执行SQL语句,DBMS_SQL.PARSE过程解析SQL语句或PL/S... 目录语法示例注意事项DBMS_SQL 是 oracle 数据库中的一个强大包,它允许动态地构建和执行

SpringBoot中使用 ThreadLocal 进行多线程上下文管理及注意事项小结

《SpringBoot中使用ThreadLocal进行多线程上下文管理及注意事项小结》本文详细介绍了ThreadLocal的原理、使用场景和示例代码,并在SpringBoot中使用ThreadLo... 目录前言技术积累1.什么是 ThreadLocal2. ThreadLocal 的原理2.1 线程隔离2

Python itertools中accumulate函数用法及使用运用详细讲解

《Pythonitertools中accumulate函数用法及使用运用详细讲解》:本文主要介绍Python的itertools库中的accumulate函数,该函数可以计算累积和或通过指定函数... 目录1.1前言:1.2定义:1.3衍生用法:1.3Leetcode的实际运用:总结 1.1前言:本文将详

浅析如何使用Swagger生成带权限控制的API文档

《浅析如何使用Swagger生成带权限控制的API文档》当涉及到权限控制时,如何生成既安全又详细的API文档就成了一个关键问题,所以这篇文章小编就来和大家好好聊聊如何用Swagger来生成带有... 目录准备工作配置 Swagger权限控制给 API 加上权限注解查看文档注意事项在咱们的开发工作里,API

Java数字转换工具类NumberUtil的使用

《Java数字转换工具类NumberUtil的使用》NumberUtil是一个功能强大的Java工具类,用于处理数字的各种操作,包括数值运算、格式化、随机数生成和数值判断,下面就来介绍一下Number... 目录一、NumberUtil类概述二、主要功能介绍1. 数值运算2. 格式化3. 数值判断4. 随机