Zeppelin 리스너 아키텍처 개선 제안 맥락 설명

차례 (13)

기존 아키텍처 및 코드 설명

Notebook과 NotebookServer의 관계

Notebook 객체는 노트를 만들고 옮기고 지우는 상위 API 입니다.

/**
 * High level api of Notebook related operations, such as create, move & delete note/folder.
 * It will also do other thing which is caused by these operation, such as update index,
 * refresh cron and update InterpreterSetting, these are done through NoteEventListener.
 */
public class Notebook {

주석의 두 번째 문장이 중요합니다. 노트 하나를 지우면 검색 색인에서 제외해야 하고, cron job 을 해제하고, 인터프리터 설정도 갱신해야 합니다. Notebook 은 이 후속 작업을 직접 하지 않고, 수신자가 누구인지 모르는 채로 NoteEventListener 에 통지합니다.

브라우저는 /ws 로 웹소켓 연결을 하나 맺습니다. 서버 쪽 끝점이 NotebookServer 입니다.

/**
 * Zeppelin websocket service. This class used setter injection because all servlet should have
 * no-parameter constructor
 */
@ManagedObject
@ServerEndpoint(value = "/ws")
public class NotebookServer implements ... {

노트를 열거나 문단을 실행하는 요청이 @OnMessage 로 수신되고, 결과는 connectionManager.broadcast() 로 열려 있는 모든 소켓에 전송됩니다. 이 경로로 다른 사용자가 변경한 내용이 실시간으로 화면에 반영됩니다.

두 객체가 서로를 필요로 하는 지점은 명확합니다. NotebookServer 는 웹소켓 메시지를 처리하려면 노트를 읽고 수정해야 하므로 Notebook 을 참조해야 합니다. Notebook 은 노트가 수정된 사실을 화면에 통지해야 하므로 리스너가 필요하고, 그 리스너가 NotebookServer 입니다.

인터프리터와 결과가 반환되는 과정

문단(paragraph)을 실행하면 구성 요소가 하나 더 관여합니다. 인터프리터는 Zeppelin 서버 프로세스가 아닌 별도 JVM 에서 실행되고, 두 프로세스는 Thrift RPC 로 통신합니다.

문단 실행 요청과 결과의 경로
NotebookServerNotebook인터프리터프로세스 경계, ThriftRemoteInterpreterProcessListener

요청은 호출 방향으로 내려갑니다. 웹소켓 메시지가 NotebookServer 에 도착하고, Notebook 을 거쳐 Paragraph 가 인터프리터로 전달됩니다. 인터프리터 프로세스가 보낸 결과 이벤트는 RemoteInterpreterEventServer 가 받아 리스너를 호출합니다.

this.listener = interpreterSettingManager.getRemoteInterpreterProcessListener();
// …
listener.onOutputUpdated(event.getNoteId(), event.getParagraphId(), event.getIndex(), …);

그 listener 가 NotebookServer 입니다. InterpreterSettingManager 의 @Inject 생성자가 RemoteInterpreterProcessListener를 받습니다.

Paragraph 상태도 같은 방식입니다. Paragraph 는 생성될 때 ParagraphJobListener 를 받아 두고, 상태가 바뀌면 그 리스너를 호출합니다.

public Paragraph(Note note, ParagraphJobListener listener) {

세 이벤트는 발생 지점이 다릅니다. 노트 변경은 Notebook 이, 문단 상태는 Paragraph 가, 인터프리터 출력은 다른 JVM 에서 전달된 것을 RemoteInterpreterEventServer 가 알립니다. 그런데 세 인터페이스의 구현체가 같은 NotebookServer 객체입니다. 브라우저에 메시지를 보낼 수 있는 곳이 웹소켓 연결을 가진 NotebookServer 하나뿐이기 때문입니다.

서로를 참조하는 두 클래스

NotebookServer 는 리스너 인터페이스 다섯 개를 한꺼번에 구현합니다.

public class NotebookServer implements AngularObjectRegistryListener,
    RemoteInterpreterProcessListener,
    ApplicationEventListener,
    ParagraphJobListener,
    NoteEventListener { ... }

HK2 바인딩도 같은 모양입니다. ZeppelinServer가 클래스 하나를 네 개의 계약으로 등록합니다.

bindAsContract(NotebookServer.class)
    .to(AngularObjectRegistryListener.class)
    .to(RemoteInterpreterProcessListener.class)
    .to(ApplicationEventListener.class)
    .to(NoteEventListener.class)
    .in(Singleton.class);

.to(NoteEventListener.class) 한 줄이 Notebook 과 NotebookServer 를 연결합니다. Notebook 은 코드 어디에서도 NotebookServer 를 직접 참조하지 않고, 생성자에 인터페이스 타입만 선언합니다.

@Inject
public Notebook(
    ZeppelinConfiguration zConf,
    // … 저장소, 권한, 인터프리터 의존성
    NoteEventListener noteEventListener)

그 파라미터에 어떤 객체가 주입될지는 위의 바인딩이 정합니다. NoteEventListener 계약에 등록된 구현이 NotebookServer 이므로, 기동할 때 HK2 가 그 싱글턴 객체를 주입합니다. Notebook.java 만 읽어서는 이 관계를 확인할 수 없습니다.

Notebook은 NotebookServer가 있어야 만들어지고, NotebookServer는 노트를 다루려면 Notebook이 있어야 합니다. 어느 쪽도 먼저 만들 수 없습니다.

노트 하나를 지우는 경로를 따라가면 이 순환이 어디서 완성되는지 확인할 수 있습니다.

노트 하나를 지울 때
removeNote리스너 호출WebSocketNotebookRestApiHTTP DELETENotebookServiceNotebookfireNoteRemoveEvent브라우저NotebookServerNoteEventListener
  • 서비스와 작업
  • 외부
  1. 삭제 요청은 REST 계층에서 서비스를 거쳐 Notebook 까지 한 방향으로 전달됩니다.
  2. Notebook 은 노트를 지운 뒤 등록된 NoteEventListener 를 차례로 호출합니다. NotebookServer 도 그중 하나이고, 통지를 받아 열려 있는 모든 웹소켓에 변경 내용을 브로드캐스트합니다.
  3. 두 객체가 서로를 참조합니다. 어느 쪽도 먼저 완성될 수 없어서 NotebookServer 는 Notebook 을 Provider 로 감싸 사용 시점에 조회합니다.같은 NotebookServer 가 웹소켓 메시지를 처리할 때는 Notebook 을 직접 호출합니다.

코드 부연 설명

Notebook 은 List<NoteEventListener> 를 필드로 보관하고, 삭제가 끝나면 차례대로 호출합니다.

private void removeNote(Note note, AuthenticationInfo subject) throws IOException {
  LOGGER.info("Remove note: {}", note.getId());
  // Set Remove to true to cancel saving this note
  note.setRemoved(true);
  noteManager.removeNote(note.getId(), subject);
  authorizationService.removeNoteAuth(note.getId());
  fireNoteRemoveEvent(note, subject);
}

반대 방향은 NotebookServer 에서 확인할 수 있습니다. 웹소켓 메시지를 처리하려면 노트를 읽어야 하기 때문에 Notebook 이 필요합니다. 하지만 필드 타입은 Notebook 이 아닙니다.

private Provider<Notebook> notebookProvider;

Provider 로 감싼 이유는 제안을 올린 메일 스레드에서 확인할 수 있습니다. 메인테이너는 이 문제를 이미 알고 있었고, HK2 를 도입할 때 NotebookServer 에 Provider 주입을 쓰기로 했다고 답했습니다. 두 객체가 서로를 요구하면 어느 쪽도 먼저 완성될 수 없어서, 한쪽의 조회를 실제 사용 시점까지 늦춰야 기동할 수 있습니다. 같은 문제를 표시한 주석이 Notebook 생성자에도 있습니다.

// TODO(zjffdu) cycle refer, not a good solution
this.interpreterSettingManager.setNotebook(this);

신규 이벤트 버스 방식 제안 및 설명

구현이 필요한 리스너 개수 분석

작년 제출한 Proposal은 코드베이스 전체의 리스너 인터페이스를 열두 항목으로 정리하고, 메서드별 호출 지점을 표로 첨부했습니다.

인터페이스비-테스트 호출 지점NotebookServer 가 구현
ApplicationEventListener16○
NoteEventListener6○
JobListener5○
RemoteInterpreterProcessListener4○
InterpreterOutputChangeListener1
AngularObjectRegistryListener0○
AngularObjectListener0
InterpreterResultMessageOutputListener0
InterpreterHookListener0

제안서의 표는 정적 분석 결과입니다. 리플렉션이나 원격 인터프리터 프로세스에서 들어오는 호출은 정적 분석에 잡히지 않아 0 으로 표시되어 있습니다. 예를 들어 AngularObjectRegistryListener 는 호출 지점이 0 이지만 NotebookServer 가 구현하고 HK2 가 바인딩합니다.

순환을 만드는 것은 NotebookServer 가 구현하는 다섯 리스너이고, 그중 호출이 잦은 것부터 옮깁니다.

피처 플래그 방식 채택

한 번에 전부 바꾸지 않습니다. 이관 도중에도 기존 동작이 유지되어야 하므로 zeppelin.eventbus.enabled 플래그로 새 경로를 켜고 끕니다.

왜 기존 의존성을 재사용하지 않았는가

zeppelin-server 는 guava 를 이미 의존성으로 가지고 있습니다. guava 의 EventBus 를 사용하면 새 의존성을 추가하지 않고 구현할 수 있습니다. 그러나 Guava 팀은 EventBus 를 사용하지 않기를 권합니다. (“We recommend against using EventBus”, EventBus 문서 참고)

오래된 설계라 더 나은 방법이 있다는 것이 이유이고, 대안으로 DI 프레임워크와 RxJava, Reactor 를 듭니다. 그중 가볍고 외부 의존이 적으며 참고할 사례가 많은 RxJava 를 채택했습니다. zeppelin-zengine 에 io.reactivex.rxjava3:rxjava:3.1.10 을 의존성으로 추가했습니다.

구현

버스 구현은 PublishSubject 한 개와 타입 필터가 전부입니다(PR 5085).

public class ZeppelinEventBus implements EventBus {
  private final Subject<Object> eventBus;

  @Inject
  public ZeppelinEventBus() {
    eventBus = PublishSubject.create();
  }

  @Override
  public void post(Object event) {
    eventBus.onNext(event);
  }

  @Override
  public <T> Observable<T> observe(Class<T> eventType) {
    return eventBus.ofType(eventType);
  }
}

플래그는 세 곳에서 확인합니다. 먼저 바인딩입니다.

if (zConf.isEventBusEnabled()) {
  bind(ZeppelinEventBus.class).to(EventBus.class).in(Singleton.class);
} else {
  bind(NoOpEventBus.class).to(EventBus.class).in(Singleton.class);
}

다음은 구독입니다.

@Inject
public void registerEventBus(EventBus eventBus, ZeppelinConfiguration zConf) {
  if (!zConf.isEventBusEnabled()) {
    LOGGER.debug("ZeppelinEventBus is disabled");
    return;
  }

  this.disposable = eventBus.observe(NoteEvent.class)
      .subscribe(this::handleNoteEvent);
}

마지막은 리스너 콜백입니다.

@Override
public void onNoteRemove(Note note, AuthenticationInfo subject) {
  if (zConf.isEventBusEnabled()) {
    return;
  }

  handleNoteRemove(note);
}

발행하는 쪽에는 조건이 없습니다. Notebook.removeNote 는 리스너를 호출하고 이벤트도 발행합니다.

fireNoteRemoveEvent(note, subject);
eventBus.post(new NoteRemoveEvent(note, subject));

두 경로가 동시에 동작하므로 같은 이벤트가 두 번 통지됩니다. 수신 측은 플래그를 확인해 리스너 콜백에서 바로 반환하고, 중복 처리를 막습니다.

이벤트 버스 도입 후ZEPPELIN-5085
removeNote발행전달WebSocketNotebookRestApiHTTP DELETENotebookServiceNotebookeventBus.post브라우저NotebookServereventBus.observeZeppelinEventBusPublishSubject
  • 서비스와 작업
  • 외부
  1. Notebook 은 구독자가 누구인지 모른 채 NoteRemoveEvent 를 버스에 발행합니다.
  2. NotebookServer 는 기동할 때 NoteEvent 타입을 구독해 두고, 전달된 이벤트로 예전 onNoteRemove 가 하던 처리를 그대로 실행합니다.

바뀌지 않는 부분

실행 방식은 바뀌지 않았습니다. 이벤트 버스라고 하면 보통 이벤트를 큐에 적재하고 즉시 반환하는 non-blocking 동작을 기대합니다. 그러나 PublishSubject.onNext 는 구독자 목록을 순회하며 같은 스레드에서 순차 호출하고, 마지막 구독자가 끝나야 반환합니다.

 // 전 — Notebook.fireNoteRemoveEvent
 for (NoteEventListener listener : noteEventListeners) {
   listener.onNoteRemove(note, subject);
 }

 // 후 — ZeppelinEventBus.post
 eventBus.onNext(event);

 // RxJava3의 onNext 구현
 public void onNext(T t) {
    for (PublishDisposable<T> pd : subscribers.get()) {
        pd.onNext(t);   // 그 자리에서 호출
    }
}

for 문이 onNext 안으로 옮겨졌을 뿐입니다. 같은 스레드에서 등록된 순서대로 호출하고, 전부 끝날 때까지 대기합니다. 바뀐 것은 Notebook 이 구독자 목록 전체를 아는 대신 이벤트 버스 하나만 안다는 점입니다(결합 방향의 역전).

post 호출이 반환되기까지
NotebookZeppelinEventBusNotebookServerremoveNotePublishSubject구독 대기

구조 그림만 보면 Notebook 이 버스에 이벤트를 발행하고 즉시 끝나는 것처럼 읽힙니다. 실제로는 구독자의 처리가 끝나야 post 가 반환합니다. 노트를 지운 스레드가 웹소켓 브로드캐스트와 잡 매니저 정리까지 끝낸 뒤에야 removeNote 에서 반환됩니다.

비동기로 처리하려면 observeOn 으로 스케줄러를 지정해야 하는데, 아직 추가하지 않았습니다.

현재 병목 지점

리뷰 단계에서 구조적 결함 발견

PR이 올라간 다음 날 CI 가 실패했습니다. NotebookTest 에서 NullPointerException 이 여러 건 발생했고, 실패한 테스트는 노트 삭제와 앵귤러 오브젝트 정리 계열에 집중되어 있었습니다.

Error:    NotebookTest.testAngularObjectRemovalOnNotebookRemove:1179 » NullPointer
Error:    NotebookTest.testAngularObjectRemovalOnParagraphRemove:1229 » NullPointer
Error:    NotebookTest.testAbortParagraphStatusOnInterpreterRestart:1461 » NullPointer

PR 작성자는 처음에 master 쪽 문제로 판단했고, 리뷰어가 CI 로그를 근거로 zengine 변경을 지목했지만, 원인은 아직 규명하지 못했습니다.

이어진 리뷰에서 리뷰어는 zeppelin-server 와 zeppelin-zengine 을 한 모듈로 합치는 것이 오래된 과제이며, 이 작업이 먼저 끝나야 한다고 지적했습니다.

새 의존성이 추가된 구조도
zeppelin-serverzeppelin-zengine호출발행전달NotebookRestApiNotebookNotebookServerZeppelinEventBusRxJava3
  • 서비스와 작업
  • 같은 시스템

이벤트 버스와 RxJava 의존성은 zeppelin-zengine 에 추가됐습니다. 그런데 구독자는 전부 zeppelin-server 쪽에 있고, zengine 의 버스를 쓰는 다른 모듈도 없습니다.

선행 작업 진행

위 리뷰에서 이어진 병합 작업이 ZEPPELIN-6355 입니다. PR 5095 가 2026-05-05 에 머지됐고, 현재 master 에는 zeppelin-zengine 디렉터리가 없습니다. Notebook.java 는 zeppelin-server/src/main/java/org/apache/zeppelin/notebook/ 아래에 있습니다.

병목으로 작용하던 조건이 사라졌습니다.

공유

관련 글

Transaction

RDBMS Transaction에 대해서 알아보겠습니다.

Query Processing

SQL이 실행 계획으로 바뀌는 과정과 디스크 I/O 비용 모델을 바탕으로 선택·외부 정렬·조인 알고리즘의 비용을 비교합니다.