惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

博客园 - 三生石上(FineUI控件)
博客园 - Franky
GbyAI
GbyAI
B
Blog
WordPress大学
WordPress大学
D
Docker
小众软件
小众软件
月光博客
月光博客
博客园 - 【当耐特】
T
The Blog of Author Tim Ferriss
IT之家
IT之家
腾讯CDC
Engineering at Meta
Engineering at Meta
Vercel News
Vercel News
H
Help Net Security
M
MIT News - Artificial intelligence
L
LangChain Blog
云风的 BLOG
云风的 BLOG
S
SegmentFault 最新的问题
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
美团技术团队
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
V
V2EX
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻

博客园 - 自由港

postgres 支持全文索引 查看 milvus 中的数据 查询数据库锁死情况 使用 keepalived 实现 tendis 高可用 部署tendis IDEA 运行 main 方法导致整个项目 install 自定义classloader hive 基础操作 在 服务器部署 seatunnel web服务 使用 seatunnel web 设计一个数据同步 如何 运行 seatunnel web 开发版 使用 seatunnel 实现数据同步 NIFI国际化 使用 NIFI读取EXCEL 数据到数据库 NIFI 使用HTTP 作为数据源接收数据 使用NIFI 同步数据库表 使用 NIFI监控数据库表 切换JDK NIFI实现配置分发
avro 数据入门
自由港 · 2025-11-09 · via 博客园 - 自由港

1.概述

Apache Avro 是一种 开源的、语言无关的、基于行的(row-based)数据序列化格式,由 Hadoop 项目开发,广泛用于大数据生态系统(如 Kafka、Spark、Flink、Hive 等)中,用于高效存储和传输结构化数据。

  • avro 是使用二进制存储
  • 需要一个schema数据格式
  • 可以通过schema 生成 类对象
  • 可以通过类对文件反序列化

2.使用方法

2.1 定义schema 文件

user.avsc

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

2.2 生成对象

定义pom.xml
plugins 定义如下

<plugins>
            <plugin>
                <groupId>org.apache.avro</groupId>
                <artifactId>avro-maven-plugin</artifactId>
                <version>1.12.0</version>
                <executions>
                    <execution>
                        <phase>generate-sources</phase>
                        <goals>
                            <goal>schema</goal>
                        </goals>
                        <configuration>
                            <sourceDirectory>${project.basedir}/src/main/resources/avro/</sourceDirectory>
                            <outputDirectory>${project.basedir}/src/main/java/</outputDirectory>
                        </configuration>
                    </execution>
                </executions>
            </plugin>
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-compiler-plugin</artifactId>
                <configuration>
                    <source>21</source>
                    <target>21</target>
                </configuration>
            </plugin>
            <plugin>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-maven-plugin</artifactId>
            </plugin>
        </plugins>

schema 文件 放到 /src/main/resources/avro/ 目录下

点击 install 后生成文件如下:

image

2.3 将对象序列化和反序列化

User user1 = new User("Alice", 256, null);
User user2 = User.newBuilder()
                .setName("Bob")
                .setFavoriteNumber(7)
                .setFavoriteColor("red")
                .build();

        // 序列化到文件
        DatumWriter<User> userDatumWriter = new SpecificDatumWriter<>(User.class);
        DataFileWriter<User> dataFileWriter = new DataFileWriter<>(userDatumWriter);
        dataFileWriter.create(user1.getSchema(), new File("users.avro"));
        dataFileWriter.append(user1);
        dataFileWriter.append(user2);
        dataFileWriter.close();

//         从文件反序列化
        DatumReader<User> userDatumReader = new SpecificDatumReader<>(User.class);
        DataFileReader<User> dataFileReader = new DataFileReader<>(new File("users.avro"), userDatumReader);
        User user = null;
        while (dataFileReader.hasNext()) {
            user = dataFileReader.next(user);
            System.out.println(user);
        }

2.4 不生成类的方式实现序列化和反序列化

每次生成对象的方式不够灵活,使用 GenericRecord 不用创建类,也可以实现AVRO数据的序列化。

Schema schema = new Schema.Parser().parse(new File("D:\\work\\research\\avrodata\\src\\main\\resources\\user.avsc"));
        GenericRecord user1 = new GenericData.Record(schema);
        user1.put("name", "Charlie");
        user1.put("favorite_number", 128);
        user1.put("favorite_color", "blue");


        DatumWriter<GenericRecord> datumWriter = new SpecificDatumWriter<>(schema);
        ByteArrayOutputStream out = new ByteArrayOutputStream();
        BinaryEncoder encoder = EncoderFactory.get().binaryEncoder(out, null);
        datumWriter.write(user1, encoder);
        encoder.flush();
        byte[] serializedBytes = out.toByteArray();

        // 从字节数组反序列化
        DatumReader<GenericRecord> datumReader = new SpecificDatumReader<>(schema);
        BinaryDecoder decoder = DecoderFactory.get().binaryDecoder(serializedBytes, null);
        GenericRecord deserializedUser = datumReader.read(null, decoder);

        System.err.println(deserializedUser);