代码之家  ›  专栏  ›  技术社区  ›  sunsets

spring kafka,MessageDeliveryException:无法将消息发送到通道

  •  0
  • sunsets  · 技术社区  · 9 年前

    让我分享每一部分,看看发生了什么。

    @EnableBinding(MultiProducerChannel.class)
    public class RealTimeDataSource{
    
    @Autowired
    RealTimeProductionService realTimeProductionService;
    
    @InboundChannelAdapter(value = MultiProducerChannel.SOURCEPRODUCTION, poller = @Poller(fixedDelay = "10000", maxMessagesPerPoll = "1"))
    public JSONArray productionMessageSource() throws Exception {
    
        long currentTime = System.currentTimeMillis();
    
        JSONArray realTimeProductionList = realTimeProductionService.getNewProductionTime();
    
        System.out.println(currentTime + " : Running source...");
        return realTimeProductionList;
    
        }
    
    }
    
    
    @EnableBinding(MultiChannel.class)
    public class RealTimeDataProcessor {
    
    @Autowired
    RealTimeProductionService realTimeProductionService;
    
    @Transformer(inputChannel = MultiChannel.PROCESSPRODUCTION, outputChannel = MultiChannel.SAVEPRODUCTION)
    public JSONObject productionMessageProcessor(List<RealTimeProduction> realTimeProductionList) throws Exception {
    
        JSONObject jsonObject = null;
        if(realTimeProductionList != null) {
            jsonObject = new JSONObject(realTimeProductionService.getNewProductionTime(realTimeProductionList));
            System.out.println("PROCESSOR RUNNING...");
        }
    
        return jsonObject;
        }
    }
    
    @EnableBinding(MultiChannel.class)
    public class RealTimeDataSink {
    
    private static final String INDEX_NAME = "c000001_kr_50879_f01";
    
    
    @Autowired
    private JestClient jestClient;
    
    @StreamListener(MultiChannel.SAVEFINALPRODUCTION)
    public void productionMessageSink(JSONObject outputs) throws Exception {
    
        if (outputs != null) {
    
            boolean indexExists = jestClient.execute(new IndicesExists.Builder(INDEX_NAME).build()).isSucceeded();
    
            JestResult jestResult = jestClient.execute(new Index.Builder(outputs).index(INDEX_NAME).type("production").build());
    
    
    
    
            System.out.println("SINK RUNNING...");
    
        }
    }
    

    }

       public interface MultiChannel {
    
    String SOURCEPRODUCTION = "production-source";
    
    String PROCESSPRODUCTION = "production-process";
    
    String SAVEPRODUCTION = "production-save";
    
    String SAVEFINALPRODUCTION = "production-save-final";
    
    @Output(SOURCEPRODUCTION)
    MessageChannel sourceproduction();
    
    @Input(PROCESSPRODUCTION)
    SubscribableChannel processorproduction();
    
    @Output(SAVEPRODUCTION)
    MessageChannel saveproduction();
    
    @Input(SAVEFINALPRODUCTION)
    SubscribableChannel savefinalproduction();
    

    }

    我认为这个资源工作得很好。但是我不知道如何找到这个错误的问题。我花了整整三天。。。。还不工作。

    org.springframework.messaging.MessageDeliveryException: failed to send Message to channel 'production-save'; nested exception is java.lang.IllegalArgumentException: payload must not be null
    

    你知道吗?

    1 回复  |  直到 9 年前