reactor-logback的AsyncAppender執(zhí)行流程源碼解讀
序
本文主要研究一下reactor-logback的AsyncAppender
AsyncAppender
reactor-logback/src/main/java/reactor/logback/AsyncAppender.java
public class AsyncAppender extends ContextAwareBase
implements Appender<ILoggingEvent>, AppenderAttachable<ILoggingEvent>,
CoreSubscriber<ILoggingEvent> {
private final AppenderAttachableImpl<ILoggingEvent> aai =
new AppenderAttachableImpl<ILoggingEvent>();
private final FilterAttachableImpl<ILoggingEvent> fai =
new FilterAttachableImpl<ILoggingEvent>();
private final AtomicReference<Appender<ILoggingEvent>> delegate =
new AtomicReference<Appender<ILoggingEvent>>();
private String name;
private WorkQueueProcessor<ILoggingEvent> processor;
private int backlog = 1024 * 1024;
private boolean includeCallerData = false;
private boolean started = false;
//......
}AsyncAppender繼承了ContextAwareBase,同時(shí)實(shí)現(xiàn)了Appender、AppenderAttachable、CoreSubscriber接口
CoreSubscriber
reactor/core/CoreSubscriber.java
public interface CoreSubscriber<T> extends Subscriber<T> {
/**
* Request a {@link Context} from dependent components which can include downstream
* operators during subscribing or a terminal {@link org.reactivestreams.Subscriber}.
*
* @return a resolved context or {@link Context#empty()}
*/
default Context currentContext(){
return Context.empty();
}
/**
* Implementors should initialize any state used by {@link #onNext(Object)} before
* calling {@link Subscription#request(long)}. Should further {@code onNext} related
* state modification occur, thread-safety will be required.
* <p>
* Note that an invalid request {@code <= 0} will not produce an onError and
* will simply be ignored or reported through a debug-enabled
* {@link reactor.util.Logger}.
*
* {@inheritDoc}
*/
@Override
void onSubscribe(Subscription s);
}CoreSubscriber繼承了Subscriber接口,Subscriber接口定義了onSubscribe(Subscription s)、onNext、onError、onComplete方法
onSubscribe
public void onSubscribe(Subscription s) {
try {
doStart();
}
catch (Throwable t) {
addError(t.getMessage(), t);
}
finally {
started = true;
s.request(Long.MAX_VALUE);
}
}
protected void doStart() {
}onSubscribe方法執(zhí)行doStart,標(biāo)記started為true,同時(shí)觸發(fā)s.request(Long.MAX_VALUE)
onNext
public void onNext(ILoggingEvent iLoggingEvent) {
aai.appendLoopOnAppenders(iLoggingEvent);
}onNext調(diào)用AppenderAttachableImpl的appendLoopOnAppenders方法
onError
public void onError(Throwable t) {
addError(t.getMessage(), t);
}onError主要是添加錯(cuò)誤信息到logback的status
onComplete
public void onComplete() {
try {
Appender<ILoggingEvent> appender = delegate.getAndSet(null);
if (appender != null){
doStop();
appender.stop();
aai.detachAndStopAllAppenders();
}
}
catch (Throwable t) {
addError(t.getMessage(), t);
}
finally {
started = false;
}
}
protected void doStop() {
}onComplete則執(zhí)行doStop、appender.stop()、aai.detachAndStopAllAppenders(),最后標(biāo)記started為false
Appender.doAppend
public void doAppend(ILoggingEvent evt) throws LogbackException {
if (getFilterChainDecision(evt) == FilterReply.DENY) {
return;
}
evt.prepareForDeferredProcessing();
if (includeCallerData) {
evt.getCallerData();
}
try {
queueLoggingEvent(evt);
}
catch (Throwable t) {
addError(t.getMessage(), t);
}
}
protected void queueLoggingEvent(ILoggingEvent evt) {
if (null != delegate.get()) {
processor.onNext(evt);
}
}doAppend方法先判斷是否需要DENY,是則直接返回,之后主要執(zhí)行queueLoggingEvent,它在delegate不為null時(shí)執(zhí)行processor.onNext(evt)
LifeCycle.start
public void start() {
startDelegateAppender();
processor = WorkQueueProcessor.<ILoggingEvent>builder().name("logger")
.bufferSize(backlog)
.autoCancel(false)
.build();
processor.subscribe(this);
}
private void startDelegateAppender() {
Appender<ILoggingEvent> delegateAppender = delegate.get();
if (null != delegateAppender && !delegateAppender.isStarted()) {
delegateAppender.start();
}
}
public void addAppender(Appender<ILoggingEvent> newAppender) {
if (delegate.compareAndSet(null, newAppender)) {
aai.addAppender(newAppender);
}
else {
throw new IllegalArgumentException(delegate.get() + " already attached.");
}
}start方法執(zhí)行startDelegateAppender,然后創(chuàng)建WorkQueueProcessor(默認(rèn)bufferSize為1024 * 1024),并subscribe當(dāng)前實(shí)例;addAppender方法會(huì)設(shè)置delegate,并往AppenderAttachableImpl添加appender
stop
public void stop() {
processor.onComplete();
}stop方法執(zhí)行processor.onComplete()
小結(jié)
reactor-logback基于WorkQueueProcessor提供了另外一種AsyncAppender,它不是基于BlockingQueue而是基于RingBuffer來實(shí)現(xiàn)的。其onSubscribe方法執(zhí)行doStart,標(biāo)記started為true,同時(shí)觸發(fā)s.request(Long.MAX_VALUE);onNext調(diào)用AppenderAttachableImpl的appendLoopOnAppenders方法;onComplete則執(zhí)行doStop、appender.stop()、aai.detachAndStopAllAppenders(),最后標(biāo)記started為false;doAppend方法先判斷是否需要DENY,是則直接返回,之后主要執(zhí)行queueLoggingEvent,它在delegate不為null時(shí)執(zhí)行processor.onNext(evt)。
以上就是reactor-logback的AsyncAppender執(zhí)行流程源碼解讀的詳細(xì)內(nèi)容,更多關(guān)于reactor-logback AsyncAppender的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!
相關(guān)文章
SpringBoot日程管理Quartz與定時(shí)任務(wù)Task實(shí)現(xiàn)詳解
定時(shí)任務(wù)是企業(yè)級(jí)開發(fā)中必不可少的組成部分,諸如長(zhǎng)周期業(yè)務(wù)數(shù)據(jù)的計(jì)算,例如年度報(bào)表,諸如系統(tǒng)臟數(shù)據(jù)的處理,再比如系統(tǒng)性能監(jiān)控報(bào)告,還有搶購類活動(dòng)的商品上架,這些都離不開定時(shí)任務(wù)。本節(jié)將介紹兩種不同的定時(shí)任務(wù)技術(shù)2022-09-09
MyBatis整合Redis實(shí)現(xiàn)二級(jí)緩存的示例代碼
這篇文章主要介紹了MyBatis整合Redis實(shí)現(xiàn)二級(jí)緩存的示例代碼,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2020-08-08
Springmvc加ajax實(shí)現(xiàn)上傳文件并頁面局部刷新
這篇文章主要介紹了Springmvc加ajax實(shí)現(xiàn)上傳文件并頁面局部刷新,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2020-06-06
mybatis簡(jiǎn)介與配置_動(dòng)力節(jié)點(diǎn)Java學(xué)院整理
這篇文章主要介紹了mybatis簡(jiǎn)介與配置,介紹了MyBatis+Spring+MySql簡(jiǎn)單配置,有興趣的可以了解一下2017-09-09

