【问题标题】:Is ManagedChannel thread-safe in Transformer KafkaTransformer Kafka 中的 ManagedChannel 线程安全吗
【发布时间】:2020-03-05 22:56:30
【问题描述】:

这是我的变压器:

public class DataEnricher implements Transformer < byte[], EnrichedData, KeyValue < byte[], EnrichedData >> {

    private ManagedChannel channel;
    private InfoClient infoclient;
    private LRUCacheCollector < String,
    InfoResponse > cache;


    public DataEnricher() {}

    @Override
    public void init(ProcessorContext context) {
        channel = ManagedChannelBuilder.forAddress("localhost", 50051).usePlaintext().build();
        infoclient = new InfoClient(channel);
    }

    @Override
    public KeyValue < byte[],
    EnrichedData > transform(byte[] key, EnrichedData request) {
        InfoResponse infoResponse = null;
        String someInfo = request.getSomeInfo();
        try {
            infoResponse = infoclient.getMoreInfo(someInfo);
        } catch (Exception e) {
            logger.warn("An exception has occurred during retrieval.", e.getMessage());
        }
        EnrichedData enrichedData = EnrichedDataBuilder.addExtraInfo(request, infoResponse);
        return new KeyValue < > (key, enrichedData);
    }

    @Override
    public KeyValue < byte[],
    DataEnricher > punctuate(long timestamp) {
        return null;
    }

    @Override
    public void close() {
        client.shutdown();
    }
}

在 Kafka Streams 中,每个流线程初始化它自己的流拓扑副本,然后根据 ProcessorContext(即每个任务,即每个分区)实例化该拓扑。那么init() 不会被调用并覆盖/泄漏每个分区的通道,并且由于我们有多个线程,甚至可以争用channel/client 的创建吗?有没有办法防止这种情况发生?

这是在run() 方法中调用的:

public KafkaStreams createStreams() {
    final Properties streamsConfiguration = new Properties();
    //other configuration is setup here
    streamsConfiguration.put(StreamsConfig.DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG, WallclockTimestampExtractor.class.getName());
    streamsConfiguration.put(
        StreamsConfig.NUM_STREAM_THREADS_CONFIG,
        3);

    StreamsBuilder streamsBuilder = new StreamsBuilder();

    RequestJsonSerde requestSerde = new RequestJsonSerde();
    DataEnricher dataEnricher = new DataEnricher();
    // Get the stream of requests
    final KStream < byte[], EnrichedData > requestsStream = streamsBuilder
        .stream(requestsTopic, Consumed.with(Serdes.ByteArray(), requestSerde));
    final KStream < byte[], EnrichedData > enrichedRequestsStream = requestsStream
        .filter((key, request) - > {
            return Objects.nonNull(request);
        }
        .transform(() - > dataEnricher);

    enrichedRequestsStream.to(enrichedRequestsTopic, Produced.with(Serdes.ByteArray()));

    return new KafkaStreams(streamsBuilder.build(), new StreamsConfig(streamsConfiguration));
}

【问题讨论】:

    标签: java multithreading apache-kafka apache-kafka-streams grpc-java


    【解决方案1】:

    ManagedChannel 无关,但您必须在TransformerSupplier 中为ProcessContext 提供新的DataEnricher 瞬间。

    KStream.transform(DataEnricher::new);
    

    一旦我遇到一些与此相关的 Kafka 流异常,我会尝试重新创建它。

    如果您不使用标点符号将更多记录发送到下游并且新密钥与输入记录相同,则 IMO 应该使用 transformValues() 因为 transform() 可能会导致重新分区,当基于密钥的操作像聚合,应用连接。

    【讨论】:

    • 谢谢,我改了,如果可行,我会通知你并接受你的回答
    【解决方案2】:

    我假设TransformerSupplier 为每个拓扑创建一个Transformer 实例(或ProcessorContext),因此每个拓扑创建一个channel。在这种情况下,channel 不会被覆盖。另外我假设你的client.shutdown() 也关闭了它的频道。

    【讨论】:

    • 我用 kafka 配置更新了问题。所以流线程数为3
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-04-15
    • 2018-07-18
    • 2011-07-04
    • 2014-04-26
    • 2012-11-30
    • 2010-12-30
    • 2013-03-12
    相关资源
    最近更新 更多