【问题标题】:Producing a MultipartFile with Spring Apache Kafka使用 Spring Apache Kafka 生成 MultipartFile
【发布时间】:2022-10-23 15:44:42
【问题描述】:

我有一个包含 String 和 MultipartFile 的模型类,我想将这个类作为 REST 方法发送到 Apache Kafka 并从那里使用它,但是我在序列化和反序列化模型类时遇到问题,我可以序列化和反序列化没有 MultipartFile 但我无法使用 MultipartFile 时的类

package com.example.demo.multipartfile.model;

import lombok.AllArgsConstructor;
import lombok.Getter;
import lombok.NoArgsConstructor;
import lombok.Setter;
import org.springframework.web.multipart.MultipartFile;

import java.io.Serializable;

@Getter
@Setter
@AllArgsConstructor
@NoArgsConstructor
public class Email implements Serializable {
    private String name;
    private MultipartFile file;

}

【问题讨论】:

  • 你得到的错误到底是什么?请不要在 Kafka 中使用 Serializeable。此外,Kafka 不适用于文件传输。因此,创建一些其他类,将您的数据解析为 JSON 等结构
  • 实际上,我们正在创建一个包含附件文件、字符串、...的服务电子邮件,我们想使用 Apache Kafka 生成和使用它。现在我重播你解决的问题,实际上我将类的所有参数转换为字节并将其作为序列化程序类中的字节数组返回并将它们拆分到反序列化程序类中,我在消费者中测试所有它们并可以再次重新创建我从 Kafka 使用的文件。如何将文件转换为 Json ,是否可行?
  • Kafka 带有 Jackson JSON 库,因此请尝试阅读 its documentation。您不应该使用 Java ByteArrayOutputStream 的原因是它非常特定于 Java,而 Kafka 客户端可以使用任何语言
  • 另外值得指出的是,Kafka 的默认最大记录大小为 1MB,电子邮件很容易大于此大小,这就是为什么不建议使用 Kafka 进行电子邮件/文件传输的原因。
  • 实际上我们可以更改请求的大小,默认为 1 MB 但它是可变的,我将其设置为 90 MB 并且它工作正常

标签: java serialization apache-kafka


【解决方案1】:

这是一个实现 MultipartFile 接口的类,当我们想要反序列化 MultipartFile,可以将字节转换为具有传输类的 MultipartFile。

package com.example.demo.multipartfile.config;

import org.springframework.core.io.Resource;
import org.springframework.web.multipart.MultipartFile;

import java.io.File;
import java.io.FileOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.nio.file.Path;

public class DecodedMultipartFile implements MultipartFile {
    private final byte [] imgContent;
    private final String originalFileName;

    public DecodedMultipartFile(byte[] imgContent , String originalFileName) {
        this.imgContent = imgContent;
        this.originalFileName = originalFileName;
    }

    @Override
    public String getName() {
        return null;
    }

    @Override
    public String getOriginalFilename() {
        return originalFileName;
    }

    @Override
    public String getContentType() {
        return null;
    }

    @Override
    public boolean isEmpty() {
        return false;
    }

    @Override
    public long getSize() {
        return 0;
    }

    @Override
    public byte[] getBytes() throws IOException {
        return imgContent;
    }

    @Override
    public InputStream getInputStream() throws IOException {
        return null;
    }

    @Override
    public Resource getResource() {
        return MultipartFile.super.getResource();
    }

    @Override
    public void transferTo(File dest) throws IOException, IllegalStateException {
        new FileOutputStream(dest).write(imgContent);

    }

    @Override
    public void transferTo(Path dest) throws IOException, IllegalStateException {
        MultipartFile.super.transferTo(dest);
    }
}

反序列化器类

package com.example.demo.multipartfile.serialization;

import com.example.demo.multipartfile.config.DecodedMultipartFile;
import com.example.demo.multipartfile.model.Email;
import lombok.Data;
import lombok.Getter;
import lombok.Setter;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.common.header.Headers;
import org.apache.kafka.common.serialization.Deserializer;
import org.springframework.web.multipart.MultipartFile;

import java.nio.ByteBuffer;
import java.util.Map;

@Slf4j
@Data
@Getter
@Setter
public class EmailDeserializer implements Deserializer<Email> {

    private String encoding = "UTF8";

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        Deserializer.super.configure(configs, isKey);
    }

    @Override
    public Email deserialize(String s, byte[] data) {
        int nameSize;
        int fileSize;
        int originalFileNameSize;

        if (data == null)
            return null;

        ByteBuffer buffer = ByteBuffer.wrap(data);

        nameSize = buffer.getInt();
        byte[] nameBytes = new byte[nameSize];
        buffer.get(nameBytes);

        fileSize = buffer.getInt();
        byte[] fileByte = new byte[fileSize];
        buffer.get(fileByte);

        originalFileNameSize = buffer.getInt();
        byte[] originalFileNameByte = new byte[originalFileNameSize];
        buffer.get(originalFileNameByte);


        try {
            String deserializedName = new String(nameBytes, encoding);
            String deserializedOriginalFileName = new String(originalFileNameByte, encoding);
            MultipartFile file = new DecodedMultipartFile(fileByte, deserializedOriginalFileName);

            return new Email(deserializedName, file);
        } catch (Exception e) {
            throw new RuntimeException(e);
        }
    }

    @Override
    public Email deserialize(String topic, Headers headers, byte[] data) {
        return Deserializer.super.deserialize(topic, headers, data);
    }

    @Override
    public void close() {
        Deserializer.super.close();
    }
}

序列化器类

package com.example.demo.multipartfile.serialization;

import com.example.demo.multipartfile.model.Email;
import lombok.Data;
import lombok.Getter;
import lombok.Setter;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.common.serialization.Serializer;

import java.io.ByteArrayOutputStream;
import java.io.DataOutputStream;
import java.io.IOException;

@Slf4j
@Data
@Getter
@Setter
public class EmailSerializer implements Serializer<Email> {
    @Override
    public byte[] serialize(String s, Email email) {

        ByteArrayOutputStream bos = new ByteArrayOutputStream();
        DataOutputStream dos = new DataOutputStream(bos);

        try {
            byte[] nameByte = email.getName().getBytes();
            byte[] fileByte = email.getFile().getBytes();
            byte[] originalFileNameByte =   email.getFile().getOriginalFilename().getBytes();

            dos.writeInt(nameByte.length);
            dos.write(nameByte);
            dos.writeInt(fileByte.length);
            dos.write(fileByte);
            dos.writeInt(originalFileNameByte.length);
            dos.write(originalFileNameByte);

            return bos.toByteArray();
        } catch (IOException e) {
            throw new RuntimeException(e);
        }
    }
}

现在我可以用 kafka 发送我想要的所有东西,序列化很奇怪,但不用担心它会起作用。 完整项目https://github.com/ehsanv8/KafkaProject.git

【讨论】:

  • 你确定这行得通吗? email.getName(),例如看起来像 return null,所以 email.getName().getBytes() 会返回一个 NPE
  • 它有点奇怪,但我测试了所有的东西,最后保存了文件,它工作正常
猜你喜欢
  • 1970-01-01
  • 2017-12-12
  • 2019-09-10
  • 1970-01-01
  • 2018-03-09
  • 2017-06-07
  • 2021-06-09
  • 2020-07-11
  • 1970-01-01
相关资源
最近更新 更多