【问题标题】:Apache Kafka create topic from code [duplicate]Apache Kafka从代码创建主题[重复]
【发布时间】:2017-04-19 02:55:16
【问题描述】:

我们知道在 Kafka 中创建Topic 应该在服务器初始化部分处理。那里我们使用默认脚本./kafka-topics --zookeeper ...,但是如果我们需要动态创建主题怎么办?

【问题讨论】:

  • 那么,如果有更好的方法,您的问题是什么?
  • @NickVanderhoven 它更像是给那些正在寻找答案的人的提示。我在文档和这里都找不到答案
  • @Andrey:请将此作为问题进行编辑,例如“如何在运行时创建 Apache Kafka 主题”,并将您的原始帖子作为新答案发布。回答你自己的帖子是可以的。
  • 嗯,够好了。无论如何:真的很棒;)我想我很快就会使用 Kafka,所以我可能会使用它。
  • @FilipMalczak 是的,我看到了混乱。将问题和答案分成两部分。

标签: java apache-kafka


【解决方案1】:

幸运的是,Kafka 0.10.1.0 为我们带来了这种能力。我在 Confluence Jira 板上看到了这些引人入胜的功能,但找不到与该主题相关的任何文档,具有讽刺意味的是,不是吗?

所以,我查看了源代码并找到了动态创建主题的方法。希望它对你们中的一些人有所帮助。当然,如果您有更好的解决方案,请不要犹豫与我们分享。

好的,让我们开始吧。

/** The method propagate topics **/
public List<String> propagateTopics(int partitions, short replication, int timeout) throws IOException {
    CreateTopicsRequest.TopicDetails topicDetails = new CreateTopicsRequest.TopicDetails(partitions, replication);
    Map<String, CreateTopicsRequest.TopicDetails> topicConfig = mTopics.stream()
            .collect(Collectors.toMap(k -> k, v -> topicDetails)); // 1

    CreateTopicsRequest request = new CreateTopicsRequest(topicConfig, timeout); // 2

    try {
        CreateTopicsResponse response = createTopic(request, BOOTSTRAP_SERVERS_CONFIG); // 3
        return response.errors().entrySet().stream()
                .filter(error -> error.getValue() == Errors.NONE)
                .map(Map.Entry::getKey)
                .collect(Collectors.toList()); // 4
    } catch (IOException e) {
        log.error(e);
    }

    return null;
}

1 我们需要一个TopicDetails 的实例,为简单起见,我将在所有主题之间共享相同的配置。假设mTopics 是您要创建的所有主题的字符串列表。

2 基本上我们想向我们的 Kafka 集群发送一个请求,现在我们有一个特殊的类,它接受 CreateTopicsRequest 和超时

3 比我们需要发送请求并获取CreateTopicsResponse

    private static final short apiKey = ApiKeys.CREATE_TOPICS.id;
    private static final short version = 0;
    private static final short correlationId = -1;

private static CreateTopicsResponse createTopic(CreateTopicsRequest request, String client) throws IllegalArgumentException, IOException {
        String[] comp = client.split(":");
        if (comp.length != 2) {
            throw new IllegalArgumentException("Wrong client directive");
        }
        String address = comp[0];
        int port = Integer.parseInt(comp[1]);

        RequestHeader header = new RequestHeader(apiKey, version, client, correlationId);
        ByteBuffer buffer = ByteBuffer.allocate(header.sizeOf() + request.sizeOf());
        header.writeTo(buffer);
        request.writeTo(buffer);

        byte byteBuf[] = buffer.array();

        byte[] resp = requestAndReceive(byteBuf, address, port);
        ByteBuffer respBuffer = ByteBuffer.wrap(resp);
        ResponseHeader.parse(respBuffer);

        return CreateTopicsResponse.parse(respBuffer);
    }

    private static byte[] requestAndReceive(byte[] buffer, String address, int port) throws IOException {
        try(Socket socket = new Socket(address, port);
            DataOutputStream dos = new DataOutputStream(socket.getOutputStream());
            DataInputStream dis = new DataInputStream(socket.getInputStream())
        ) {
            dos.writeInt(buffer.length);
            dos.write(buffer);
            dos.flush();

            byte resp[] = new byte[dis.readInt()];
            dis.readFully(resp);

            return resp;
        } catch (IOException e) {
            log.error(e);
        }

        return new byte[0];
    }

这里根本没有魔法,只是发送请求,然后将字节流解析为响应。

4 CreateTopicsResponse 具有属性errors,它只是一个Map&lt;String, Errors&gt;,其中key 是您请求的主题名称。棘手的是,它包含您请求的所有主题,但没有错误的主题具有值Errors.None,这就是我过滤响应并仅返回成功创建的主题的原因。

【讨论】:

  • 对 Kafka 0.10.2.0 使用 version=1
【解决方案2】:

扩展 Andrei Nechaev 的答案

在 10.2.0 中,获取 CreateTopicsRequest 实例的方式发生了一些变化。我们需要使用 Builder 内部类来构建一个 CreateTopicsRequest 实例。这是一个代码示例。

CreateTopicsRequest.Builder builder = new CreateTopicsRequest.Builder(topicConfig, timeout, false);
CreateTopicsRequest request = builder.build();

【讨论】:

    猜你喜欢
    • 2018-03-31
    • 2016-07-26
    • 2016-07-21
    • 2017-07-21
    • 1970-01-01
    • 2021-09-09
    • 2016-02-05
    • 2020-01-22
    相关资源
    最近更新 更多