You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

求基于Kafka Connect SMT实现Claim-Check Pattern的实例与参考

使用Kafka Connect SMT实现Claim-Check模式(Azure环境)

针对你在Azure Blob Storage + Azure Event Hub环境下的Claim-Check模式需求,以下是实际落地的配置示例、自定义SMT逻辑和关键要点:

一、核心实现思路

  1. Sink端导出:通过Sink连接器+SMT将原消息中的大负载上传至Azure Blob Storage,将原消息替换为包含Blob访问URL的元数据后发送到Azure Event Hub
  2. Source端导入:Source连接器从Event Hub读取元数据消息,通过SMT根据Blob URL拉取负载,还原完整消息后传递给下游处理

二、Sink端配置(导出大负载到Blob)

自定义AzureClaimCheckOut SMT核心逻辑

若官方SMT未适配Azure Blob,可基于Azure Storage SDK自定义SMT实现上传逻辑:

public class AzureClaimCheckOut<S extends ConnectRecord<S>> extends Transformation<S> {
    private BlobContainerClient containerClient;

    @Override
    public void configure(Map<String, ?> configs) {
        // 初始化Azure Blob容器客户端
        String accountName = (String) configs.get("storage.azure.account.name");
        String accountKey = (String) configs.get("storage.azure.account.key");
        String containerName = (String) configs.get("storage.container");
        StorageSharedKeyCredential credential = new StorageSharedKeyCredential(accountName, accountKey);
        String endpoint = String.format(Locale.ROOT, "https://%s.blob.core.windows.net", accountName);
        BlobServiceClient serviceClient = new BlobServiceClientBuilder().endpoint(endpoint).credential(credential).buildClient();
        containerClient = serviceClient.getBlobContainerClient(containerName);
    }

    @Override
    public S apply(S record) {
        // 提取大负载字段
        Struct value = (Struct) record.value();
        byte[] payload = value.getBytes("payload");
        // 生成唯一Blob名称(用消息offset+时间戳避免重复)
        String blobName = String.format("msg-%d-%d", record.kafkaOffset(), System.currentTimeMillis());
        // 上传到Blob
        BlobClient blobClient = containerClient.getBlobClient(blobName);
        blobClient.upload(BinaryData.fromBytes(payload), payload.length);
        // 替换消息为元数据结构
        Struct newValue = new Struct(record.valueSchema());
        newValue.put("blob_url", blobClient.getBlobUrl());
        newValue.put("metadata", value.get("metadata")); // 保留原消息的业务元数据
        return record.newRecord(record.topic(), record.kafkaPartition(), record.keySchema(), record.key(), newValue.schema(), newValue, record.timestamp());
    }

    // 省略其他接口实现代码
}

连接器配置示例

name=eventhub-sink-claimcheck
connector.class=com.microsoft.azure.eventhubs.kafka.connect.sink.EventHubSinkConnector
tasks.max=3
topics=large-messages-topic
eventhubs.namespace=your-eventhub-namespace
eventhubs.name=your-eventhub-name
eventhubs.sas.key=your-sas-key
eventhubs.sas.key.name=your-sas-key-name

# 自定义SMT配置
transforms=azureClaimCheckOut
transforms.azureClaimCheckOut.type=com.yourcompany.connect.transforms.AzureClaimCheckOut
transforms.azureClaimCheckOut.storage.azure.account.name=your-storage-account
transforms.azureClaimCheckOut.storage.azure.account.key=your-storage-key
transforms.azureClaimCheckOut.storage.container=claimcheck-blobs

三、Source端配置(从Blob加载负载)

自定义AzureClaimCheckIn SMT核心逻辑

public class AzureClaimCheckIn<S extends ConnectRecord<S>> extends Transformation<S> {
    private BlobServiceClient serviceClient;

    @Override
    public void configure(Map<String, ?> configs) {
        // 初始化Azure Blob服务客户端(同Sink端逻辑)
        String accountName = (String) configs.get("storage.azure.account.name");
        String accountKey = (String) configs.get("storage.azure.account.key");
        StorageSharedKeyCredential credential = new StorageSharedKeyCredential(accountName, accountKey);
        String endpoint = String.format(Locale.ROOT, "https://%s.blob.core.windows.net", accountName);
        serviceClient = new BlobServiceClientBuilder().endpoint(endpoint).credential(credential).buildClient();
    }

    @Override
    public S apply(S record) {
        // 读取元数据中的Blob URL
        Struct value = (Struct) record.value();
        String blobUrl = value.getString("blob_url");
        // 解析Blob名称并下载内容
        URI blobUri = URI.create(blobUrl);
        String blobName = blobUri.getPath().substring(1);
        BlobClient blobClient = serviceClient.getBlobContainerClient(blobUri.getHost().split("\\.")[0]).getBlobClient(blobName);
        BinaryData blobData = blobClient.downloadContent();
        byte[] payload = blobData.toBytes();
        // 还原完整消息结构
        Struct newValue = new Struct(record.valueSchema());
        newValue.put("payload", payload);
        newValue.put("metadata", value.get("metadata"));
        return record.newRecord(record.topic(), record.kafkaPartition(), record.keySchema(), record.key(), newValue.schema(), newValue, record.timestamp());
    }

    // 省略其他接口实现代码
}

连接器配置示例

name=eventhub-source-claimcheck
connector.class=com.microsoft.azure.eventhubs.kafka.connect.source.EventHubSourceConnector
tasks.max=3
eventhubs.namespace=your-eventhub-namespace
eventhubs.name=your-eventhub-name
eventhubs.sas.key=your-sas-key
eventhubs.sas.key.name=your-sas-key-name
topic=processed-messages-topic

# 自定义SMT配置
transforms=azureClaimCheckIn
transforms.azureClaimCheckIn.type=com.yourcompany.connect.transforms.AzureClaimCheckIn
transforms.azureClaimCheckIn.storage.azure.account.name=your-storage-account
transforms.azureClaimCheckIn.storage.azure.account.key=your-storage-key

四、落地关键注意事项

  • 权限优化:避免硬编码存储账户密钥,改用Azure AD服务主体认证,为连接器分配Storage Blob Data Contributor(写入)和Storage Blob Data Reader(读取)角色
  • 幂等性保障:Blob命名使用消息的kafkaOffset+UUID组合,防止重复上传覆盖;开启连接器幂等性配置,避免重复消费
  • 错误处理:配置连接器死信队列(DLQ),将Blob下载失败、消息结构异常的消息路由到指定主题,便于后续排查
  • 性能调优:针对超大文件使用Blob分块上传/下载API;调整连接器batch.size参数,平衡吞吐量和内存占用

内容的提问来源于stack exchange,提问作者gt1

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.15 17:43:19