【问题标题】:Cloud stream not able to track the status for down stream failures云流无法跟踪下游故障的状态
【发布时间】:2022-10-24 11:50:44
【问题描述】:

我编写了以下代码来利用云流functional approach 从 RabbitMQ 获取事件并将其发布到 KAFKA,如果 KAFKA 代理由于任何原因而停止运行,我可以在运行应用程序时实现主要目标原因然后我得到了 KAFKA BROKER 的日志,它已关闭,但同时我想停止来自 rabbitMQ 的事件,或者直到代理出现之前,这些消息应该被路由到 ExchangeDLQ topic。但是,我在很多地方都看到过使用producer sync: true,但在我的情况下,这要么没有帮助,很多人在目标频道出现故障时谈论@ServiceActivator(inputChannel = "error-topic") for error topic,这种方法是也没有被执行。所以简而言之,我不想丢失在kafka期间从rabbitMQ收到的消息由于任何原因而关闭

应用程序.yml

management:
  health:
    binders:
      enabled: true
    kafka:
      enabled: true
server:
  port: 8081

spring:
  rabbitmq:
    publisher-confirms : true
  kafka:
    bootstrap-servers: localhost:9092
    producer:
      properties:
        max.block.ms: 100
    admin:
      fail-fast: true
  cloud:
    function:
      definition: handle
    stream:
      bindingRetryInterval : 30
      rabbit:
        bindings:
          handle-in-0:
            consumer:
              bindingRoutingKey: MyRoutingKey
              exchangeType: topic
              requeueRejected : true
              acknowledgeMode: AUTO
      #              ackMode: MANUAL
      #              acknowledge-mode: MANUAL
      #              republishToDlq : false
      kafka:
        binder:
          considerDownWhenAnyPartitionHasNoLeader: true
          producer:
            properties:
              max.block.ms : 100
          brokers:
            - localhost
      bindings:
        handle-in-0:
          destination: test_queue
          binder: rabbit
          group: queue
        handle-out-0:
          destination: mytopic
          producer:
            sync: true
            errorChannelEnabled: true
          binder: kafka
      binders:
        error:
          destination: myerror
        rabbit:
          type: rabbit
          environment:
            spring:
              rabbitmq:
                host: localhost
                port: 5672
                username: guest
                password: guest
                virtual-host: rahul_host
        kafka:
          type: kafka


json:
  cuttoff:
    size:
      limit: 1000

CloudStreamConfig.java

@Configuration
public class CloudStreamConfig {
    private static final Logger log = LoggerFactory.getLogger(CloudStreamConfig.class);

    @Autowired
    ChunkService chunkService;

    @Bean
    public Function<Message<RmaValues>,Collection<Message<RmaValues>>> handle() {
        return rmaValue -> {
            log.info("processor runs : message received with request id : {}", rmaValue.getPayload().getRequestId());
            ArrayList<Message<RmaValues>> msgList = new ArrayList<Message<RmaValues>>();
            try {
                List<RmaValues> dividedJson = chunkService.getDividedJson(rmaValue.getPayload());
                for(RmaValues rmaValues : dividedJson) {
                    msgList.add(MessageBuilder.withPayload(rmaValues).build());
                }
            } catch (Exception e) {
                e.printStackTrace();

            }
            Channel channel = rmaValue.getHeaders().get(AmqpHeaders.CHANNEL, Channel.class);
            Long deliveryTag = rmaValue.getHeaders().get(AmqpHeaders.DELIVERY_TAG, Long.class);

//            try {
//                channel.basicAck(deliveryTag, false);
//            } catch (IOException e) {
//                e.printStackTrace();
//            }
            return msgList;
        };
    };
    @ServiceActivator(inputChannel = "error-topic")
    public void errorHandler(ErrorMessage em) {
        log.info("---------------------------------------got error message over errorChannel: {}", em);
        if (null != em.getPayload() && em.getPayload() instanceof KafkaSendFailureException) {
            KafkaSendFailureException kafkaSendFailureException = (KafkaSendFailureException) em.getPayload();
            if (kafkaSendFailureException.getRecord() != null && kafkaSendFailureException.getRecord().value() != null
                    && kafkaSendFailureException.getRecord().value() instanceof byte[]) {
                log.warn("error channel message. Payload {}", new String((byte[])(kafkaSendFailureException.getRecord().value())));
            }
        }
    }

KafkaProducerConfiguration.java

@Configuration
        public class KafkaProducerConfiguration {
        
            @Value(value = "${spring.kafka.bootstrap-servers}")
            private String bootstrapAddress;
        
            @Bean
            public ProducerFactory<String, Object> producerFactory() {
                Map<String, Object> configProps = new HashMap<>();
                configProps.put(
                        ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
                        bootstrapAddress);
                configProps.put(
                        ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
                        StringSerializer.class);
                configProps.put(
                        ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
                        StringSerializer.class);
                return new DefaultKafkaProducerFactory<>(configProps);
            }
        
            @Bean
            public KafkaTemplate<String, String> kafkaTemplate() {
                return new KafkaTemplate(producerFactory());
            }
        

RmModelOutputIngestionApplication.java

@SpringBootApplication(scanBasePackages = "com.abb.rm")
        public class RmModelOutputIngestionApplication {
            private static final Logger LOGGER = LogManager.getLogger(RmModelOutputIngestionApplication.class);
        
            public static void main(String[] args) {
                SpringApplication.run(RmModelOutputIngestionApplication.class, args);
            }
        
            @Bean("objectMapper")
            public ObjectMapper objectMapper() {
                ObjectMapper mapper = new ObjectMapper();
                LOGGER.info("Returning object mapper...");
            return mapper;
        }

【问题讨论】:

    标签: spring functional-programming cloud spring-cloud-stream spring-rabbit


    【解决方案1】:

    首先,您似乎创建了太多不必要的代码。为什么你有 ObjectMapper?为什么你有 KafkaTemplate?为什么有ProducerFactory?这些都已经为您提供了。 你真的只需要一个函数,可能还有一个错误处理程序——取决于你选择的错误处理策略,这让我想到了错误处理主题。处理错误的主要方法有 3 种。 Here is the link 向文档解释所有内容并提供样本。请通读并相应地修改您的应用程序,如果某些内容不起作用或不清楚,请随时跟进。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-01-29
      • 2013-12-08
      • 1970-01-01
      • 1970-01-01
      • 2012-01-28
      • 2014-06-27
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多