【问题标题】:Unable to insert data into db via an Observable无法通过 Observable 将数据插入数据库
【发布时间】:2020-07-27 14:26:02
【问题描述】:

我有以下 Observable,我希望在订阅它时会发生一些 DB 插入。 但是什么也没发生,没有数据库插入,同时也没有错误。

但如果我直接订阅执行 DB 调用的方法,则 DB 插入会按预期发生。 我该如何解决这个问题,以便订阅下面的 Observable 调用将执行数据库插入?

请指教。谢谢。

这是没有发生数据库插入且没有错误的 Observable。我想改变这一点,以便在订阅此 Observable 时发生 DB 插入。

public Observable<KafkaConsumerRecord<String, RequestObj>> apply(KafkaConsumerRecords<String, RequestObj> records) {

    Observable.from(records.getDelegate().records().records("TOPIC_NAME"))
            .buffer(2)
            .map(this::convertToEventRequest)
            .doOnNext(this::handleEventInsertions)
            .doOnSubscribe(() -> System.out.println("Subscribed!"))
            .subscribe(); // purposely subscribing here itself to test 

    return null; // even if I return this observable and subscribe at the caller, same outcome. 
}

只是为了测试查询是否有效,如果我要直接订阅执行插入的方法,它会按预期工作,如下所示。 在调试模式下执行此操作。

client.rxQueryWithParams(query, new JsonArray(params)).subscribe() // works

以下是查看 convertToEventRequest 和 handleEventInsertions 方法内部发生的事情的参考

private Map<String, List<?>> convertToEventRequest(Object records) {
    List<ConsumerRecord<String, RequestObj>> consumerRecords = (List<ConsumerRecord<String, RequestObj>>) records;

    List<AddEventRequest> addEventRequests = new ArrayList<>();
    List<UpdateEventRequest> updateEventRequests = new ArrayList<>();

    consumerRecords.forEach(record -> {
        String eventType = new String(record.headers().headers("type").iterator().next().value(), StandardCharsets.UTF_8);

        if("add".equals(eventType)) {
            AddEventRequest request = AddEventRequest.builder()
                    .count(Integer.parseInt(new String(record.headers().headers("count").iterator().next().value(), StandardCharsets.UTF_8)))
                    .data(record.value())
                    .build();
            addEventRequests.add(request);
        } else {
            UpdateEventRequest request = UpdateEventRequest.builder()
                    .id(new String(record.headers().headers("id").iterator().next().value(), StandardCharsets.UTF_8))
                    .status(Integer.parseInt(new String(record.headers().headers("status").iterator().next().value(), StandardCharsets.UTF_8)))
                    .build();
            updateEventRequests.add(request);
        }
    });

    return new HashMap<String, List<?>>() {{
        put("add", addEventRequests);
        put("update", updateEventRequests);
    }};
}

private void handleEventInsertions(Object eventObject) {
    Map<String, List<?>> eventMap = (Map<String, List<?>>) eventObject;

    List<AddEventRequest> addEventRequests = (List<AddEventRequest>) eventMap.get("add");
    List<UpdateEventRequest> updateEventRequests = (List<UpdateEventRequest>) eventMap.get("update");

    if(addEventRequests != null && !addEventRequests.isEmpty()) {
        insertAddEvents(addEventRequests);
    }
    if(updateEventRequests != null && !updateEventRequests.isEmpty()) {
        insertUpdateEvents(updateEventRequests);
    }
}

private Single<ResultSet> insertAddEvents(List<AddEventRequest> requests) {
    AddEventRequest request = requests.get(0);
    List<Object> params = Arrays.asList(request.getCount(), request.getData());
    String query = "INSERT INTO mytable(count, data, creat_ts) " +
            "VALUES (?, ?, current_timestamp)";
    return client.rxQueryWithParams(query, new JsonArray(params));
}

private Single<ResultSet> insertUpdateEvents(List<UpdateEventRequest> requests) {
    UpdateEventRequest request = requests.get(0);
    return client.rxQueryWithParams(
            "UPDATE mytable SET status=?, creat_ts=current_timestamp WHERE id=?",
            new JsonArray(Arrays.asList(request.getStatus(), request.getId())));
}

【问题讨论】:

    标签: java observable rx-java


    【解决方案1】:

    你能试着把它包装成Observable.defer吗?

    Observable.defer(() -> Observable.from(records.getDelegate().records().records("TOPIC_NAME"))...
    

    【讨论】:

    • 如上所述将其包装到 defer 中没有任何区别。
    猜你喜欢
    • 1970-01-01
    • 2017-06-21
    • 2012-04-13
    • 1970-01-01
    • 2014-01-13
    • 1970-01-01
    相关资源
    最近更新 更多