Merge pull request #883 from Netyyyy/myao/changerfdoc

Change subscribe() to block()
This commit is contained in:
Muyao Feng 2022-03-24 16:02:03 +08:00 коммит произвёл GitHub
Родитель 44bc6e20b8 59d03ff863
Коммит 8f3601fb08
Не найден ключ, соответствующий данной подписи
Идентификатор ключа GPG: 4AEE18F83AFDEB23
2 изменённых файлов: 10 добавлений и 10 удалений

Просмотреть файл

@ -363,7 +363,7 @@ public Consumer<Message<String>> consume() {
checkpointer.success()
.doOnSuccess(success -> LOGGER.info("Message '{}' successfully checkpointed", message.getPayload()))
.doOnError(error -> LOGGER.error("Exception found", error))
.subscribe();
.block();
};
}
@ -456,7 +456,7 @@ public Consumer<Message<List<String>>> consume() {
checkpointer.success()
.doOnSuccess(success -> LOGGER.info("Message '{}' successfully checkpointed", message.getPayload()))
.doOnError(error -> LOGGER.error("Exception found", error))
.subscribe();
.block();
};
}
@ -612,7 +612,7 @@ public Consumer<Message<String>> consume1() {
checkpointer.success()
.doOnSuccess(success -> LOGGER.info("Message1 '{}' successfully checkpointed", message))
.doOnError(error -> LOGGER.error("Exception found", error))
.subscribe();
.block();
};
}
@ -624,7 +624,7 @@ public Consumer<Message<String>> consume2() {
checkpointer.success()
.doOnSuccess(success -> LOGGER.info("Message2 '{}' successfully checkpointed", message))
.doOnError(error -> LOGGER.error("Exception found", error))
.subscribe();
.block();
};
}
----
@ -920,7 +920,7 @@ public Consumer<Message<String>> consume() {
checkpointer.success()
.doOnSuccess(success -> LOGGER.info("Message '{}' successfully checkpointed", message.getPayload()))
.doOnError(error -> LOGGER.error("Exception found", error))
.subscribe();
.block();
};
}
@ -1112,7 +1112,7 @@ public Consumer<Message<String>> consume1() {
checkpointer.success()
.doOnSuccess(s -> LOGGER.info("Message '{}' successfully checkpointed", message.getPayload()))
.doOnError(e -> LOGGER.error("Error found", e))
.subscribe();
.block();
};
}
@ -1124,7 +1124,7 @@ public Consumer<Message<String>> consume2() {
checkpointer.success()
.doOnSuccess(s -> LOGGER.info("Message '{}' successfully checkpointed", message.getPayload()))
.doOnError(e -> LOGGER.error("Error found", e))
.subscribe();
.block();
};
}

Просмотреть файл

@ -287,7 +287,7 @@ class Demo {
checkpointer.success()
.doOnSuccess(s -> LOGGER.info("Message '{}' successfully checkpointed", message))
.doOnError(e -> LOGGER.error("Error found", e))
.subscribe();
.block();
}
}
----
@ -634,7 +634,7 @@ class Demo {
checkpointer.success()
.doOnSuccess(s -> LOGGER.info("Message '{}' successfully checkpointed", message))
.doOnError(e -> LOGGER.error("Error found", e))
.subscribe();
.block();
}
}
----
@ -960,7 +960,7 @@ class Demo {
checkpointer.success()
.doOnError(Throwable::printStackTrace)
.doOnSuccess(t -> LOGGER.info("Message '{}' successfully checkpointed", message))
.subscribe();
.block();
}
}
----