확장
확장 (Extensions) 만들기
Druid는 런타임에 확장(extension)을 추가할 수 있게 해 주는 모듈 시스템을 사용해요.
출처: 문서
본문
자신만의 확장 작성하기
Druid의 확장은 Guice를 활용해 런타임에 무언가를 추가합니다. 기본적으로 Guice는 의존성 주입(Dependency Injection) 프레임워크인데, Druid는 이를 사용해서 Druid 프로세스의 예상 객체 그래프를 담아 둡니다. 확장은 Guice 바인딩(binding)을 추가해서 객체 그래프에 원하는(필요한) 변경을 가할 수 있어요. 확장은 사실상 거의 무엇이든 원하는 대로 바꿀 수 있는 능력을 주지만, 일반적으로 사람들은 아래 나열된 것 중 하나를 확장하기를 원한다고 봅니다. 이는 이 페이지에서 언급하는 인터페이스에 영향을 주는 변경에 대해 우리의 버전 관리 전략(versioning strategy)을 지킨다는 뜻이고, 다른 인터페이스는 "internal"로 간주되어 patch 릴리스 사이에도 호환되지 않는 방식으로 변경될 수 있습니다.
org.apache.druid.segment.loading.DataSegment*과org.apache.druid.tasklogs.TaskLog*클래스를 확장해서 새 딥 스토리지 구현 추가org.apache.druid.data.input.InputSource를 확장해서 새 input source 추가org.apache.druid.data.input.InputEntity를 확장해서 새 input entity 추가- 필요하다면
org.apache.druid.data.input.InputSourceReader를 확장해서 새 input source reader 추가. 대부분의 경우org.apache.druid.data.input.impl.InputEntityIteratingReader를 사용할 수 있음 org.apache.druid.data.input.InputFormat을 확장해서 새 input format 추가- 텍스트 형식이면
org.apache.druid.data.input.TextReader, 바이너리 형식이면org.apache.druid.data.input.IntermediateRowParsingReader를 확장해서 새 input entity reader 추가 org.apache.druid.query.aggregation.AggregatorFactory,org.apache.druid.query.aggregation.Aggregator,org.apache.druid.query.aggregation.BufferAggregator를 확장해서 Aggregator 추가org.apache.druid.query.aggregation.PostAggregator를 확장해서 PostAggregator 추가org.apache.druid.query.extraction.ExtractionFn을 확장해서 ExtractionFn 추가org.apache.druid.segment.serde.ComplexMetricSerde를 확장해서 Complex metric 추가org.apache.druid.query.QueryRunnerFactory,org.apache.druid.query.QueryToolChest,org.apache.druid.query.Query를 확장해서 새 Query 타입 추가Jerseys.addResource(binder, clazz)를 호출해서 새 Jersey 리소스 추가org.apache.druid.server.initialization.jetty.ServletFilterHolder를 확장해서 새 Jetty 필터 추가org.apache.druid.metadata.PasswordProvider를 확장해서 새 시크릿 제공자(secret provider) 추가org.apache.druid.metadata.DynamicConfigProvider를 확장해서 새 동적 구성 제공자(dynamic config provider) 추가druid-processing패키지의org.apache.druid.segment.transform.Transform인터페이스를 구현해서 새 수집 transform 추가- 확장을 다른 모든 Druid 확장과 함께 번들링
확장은 org.apache.druid.initialization.DruidModule의 구현을 통해 시스템에 추가됩니다.
Druid Module 만들기
DruidModule 클래스에는 두 가지 메서드가 있어요.
configure(Binder)메서드getJacksonModules()메서드
configure(Binder) 메서드는 일반 Guice 모듈이 갖는 메서드와 동일합니다. getJacksonModules() 메서드는 Druid가 사용하는 Jackson ObjectMapper 인스턴스 초기화를 돕는 Jackson 모듈 목록을 제공해요. 이것이 Jackson을 통해 인스턴스화되는 확장(AggregatorFactory, InputSource 객체 같은)을 Druid에 추가하는 방법입니다.
Druid Module 등록하기
DruidModule을 만들었다면 jar의 META-INF/services 디렉토리에 추가 파일을 패키징해야 해요. 이는 maven 프로젝트에서 src/main/resources 디렉토리에 파일을 만들어서 가장 쉽게 할 수 있습니다. 이에 대한 예시는 Druid 코드의 cassandra-storage, hdfs-storage, s3-extensions 모듈에 있습니다.
jar에 존재해야 하는 파일은
META-INF/services/org.apache.druid.initialization.DruidModule
이며, DruidModule을 구현하는 패키지-완전 클래스 이름들을 개행으로 구분해 나열한 텍스트 파일이어야 합니다. 예:
org.apache.druid.storage.cassandra.CassandraDruidModule
jar에 이 파일이 있으면, 클래스패스에 추가되거나 확장으로 추가될 때 Druid가 그 파일을 발견하고 Module 인스턴스들을 인스턴스화합니다. Module은 기본 생성자를 가져야 하지만, 런타임 구성 속성에 접근해야 한다면 Properties 객체가 Guice에서 주입되도록 하는 @Inject 어노테이션을 가진 메서드를 가질 수 있어요.
새 딥 스토리지 구현 추가
이를 수행하는 예시는 druid-azure-extensions, druid-google-extensions, druid-cassandra-storage, druid-hdfs-storage, druid-s3-extensions 모듈을 확인하세요.
확장의 기본 아이디어는 DataSegmentPusher와 URIDataPuller 객체에 대한 바인딩을 추가해야 한다는 것입니다. 추가하는 방법은 다음과 같습니다(HdfsStorageDruidModule에서 가져옴):
Binders.dataSegmentPullerBinder(binder) .addBinding("hdfs") .to(HdfsDataSegmentPuller.class).in(LazySingleton.class);Binders.dataSegmentPusherBinder(binder) .addBinding("hdfs") .to(HdfsDataSegmentPusher.class).in(LazySingleton.class);
Binders.dataSegment*Binder()는 Guice multibind "MapBinder"를 설정해 주는 druid-core jar가 제공하는 호출이에요. 그 점이 이해되지 않아도 걱정하지 마세요. 마법 같은 주문이라고 생각하면 됩니다.
Puller 바인더의 addBinding("hdfs")는 "hdfs" 타입의 loadSpec 객체에 대한 새 핸들러를 만들고, Pusher 바인더의 경우 druid.storage.type 파라미터에 지정할 수 있는 새 타입 값을 만들어요.
to(...).in(...);은 일반적인 Guice 코드입니다.
DataSegmentPusher와 URIDataPuller 외에도 다음을 바인딩할 수 있어요.
DataSegmentKiller: 세그먼트를 제거함. Kill Task의 일부로 사용되어 사용하지 않는 세그먼트를 삭제하는, 즉 더 새로운 버전에 대체되었거나 클러스터에서 드롭된 세그먼트들의 가비지 컬렉션을 수행DataSegmentMover: 세그먼트를 한 곳에서 다른 곳으로 마이그레이션하도록 허용. 현재는 MoveTask의 일부로 사용하지 않는 세그먼트를 다른 S3 버킷·접두사로 옮기는 데만 사용. 주로 사용하지 않는 데이터의 스토리지 비용을 줄이기 위해(예: glacier나 더 싼 스토리지로 이동)DataSegmentArchiver: Mover의 래퍼이지만 미리 구성된 대상 버킷/경로가 함께 와서, ArchiveTask의 일부로 런타임에 지정할 필요가 없음
딥 스토리지 구현 검증하기
경고! 이는 공식적인 절차가 아니라, 새 딥 스토리지 구현이 push·pull·kill 세그먼트를 할 수 있는지 검증하는 힌트 모음입니다.
구현 검증에는 배치 수집 태스크를 사용하는 것을 권장해요. 세그먼트는 약 1분 후에 자동으로 Historical 노드로 롤업됩니다. 이런 방식으로 push(realtime 프로세스에서)와 pull(Historical 프로세스에서) 세그먼트 양쪽을 검증할 수 있습니다.
DataSegmentPusher
데이터 스토리지(클라우드 스토리지 서비스, 분산 파일 시스템 등)가 어디든, 수집 태스크가 끝난 후 index.zip(HDFS 데이터 스토리지면 partitionNum_index.zip) 파일 하나가 새로 생긴 것을 볼 수 있어야 해요.
URIDataPuller
수집 태스크가 끝난 후 약 1분 뒤, Historical 프로세스가 새 세그먼트를 로드하려고 시도하는 것을 볼 수 있어야 합니다.
DataSegmentKiller
세그먼트 제거를 테스트하는 가장 쉬운 방법은 세그먼트를 사용하지 않는 것으로 표시한 뒤 웹 콘솔에서 kill 태스크를 시작하는 것이에요.
세그먼트를 사용하지 않는 것으로 표시하려면 메타데이터 스토리지에 연결해서 세그먼트 테이블 행의 used 컬럼을 false로 업데이트해야 합니다. 세그먼트 kill 태스크를 시작하려면 웹 콘솔에 접속한 뒤 적절한 datasource에 대해 issue kill task를 선택하면 됩니다. kill 태스크가 끝난 후 index.zip(HDFS 데이터 스토리지면 partitionNum_index.zip) 파일이 데이터 스토리지에서 삭제돼야 합니다.
새 input source 지원 추가
새 input source 지원을 추가하려면 InputSource, InputEntity, InputSourceReader 세 개의 인터페이스를 구현해야 해요. InputSource는 입력 데이터가 어디에 저장되는지 정의하고, InputEntity는 네이티브 병렬 인덱싱에서 데이터를 어떻게 병렬로 읽을 수 있는지 정의하며, InputSourceReader는 새 input source를 어떻게 읽을지 정의합니다. 대부분의 경우 제공되는 InputEntityIteratingReader를 간단히 사용할 수 있어요.
이에 대한 예시는 druid-s3-extensions 모듈에 S3InputSource와 S3Entity로 있습니다.
InputSource 추가는 Guice 대신 거의 전적으로 Jackson Modules를 통해 이루어져요. 구체적으로 구현을 주목하세요:
@Overridepublic List<? extends Module> getJacksonModules(){ return ImmutableList.of( new SimpleModule().registerSubtypes(new NamedType(S3InputSource.class, "s3")) );}
이것은 InputSource를 Jackson의 다형성(polymorphic) 직렬화/역직렬화 레이어에 등록하는 것입니다. 더 구체적으로 말하면, IO config에 "inputSource": { "type": "s3", ... }를 지정하면 시스템이 당신의 InputSource 구현을 위한 이 InputSource를 로드한다는 뜻이에요.
참고로 Druid 안에서는 Jackson 역직렬화 객체에 대한 @JacksonInject 어노테이션이 실제로 기본 Guice 인젝터를 사용해 주입할 객체를 해석하도록 만들었어요. 그래서 InputSource가 어떤 객체에 접근해야 한다면 setter에 @JacksonInject 어노테이션을 추가하면 인스턴스화 시점에 설정됩니다.
새 데이터 형식 지원 추가
새 데이터 형식 지원을 추가하려면 InputFormat과 InputEntityReader 두 인터페이스를 구현해야 해요. InputFormat은 데이터가 어떻게 형식화되는지 정의하고, InputEntityReader는 데이터를 어떻게 파싱해서 Druid InputRow로 변환하는지 정의합니다.
druid-orc-extensions 모듈에 OrcInputFormat과 OrcReader로 예시가 있습니다.
InputFormat 추가는 InputSource 추가와 매우 비슷해요. 둘 다 순전히 Jackson을 통해 동작하므로 단지 DruidModule이 반환하는 Jackson 모듈에 추가만 하면 됩니다.
Aggregator 추가
AggregatorFactory 객체 추가는 InputSource 객체와 매우 비슷합니다. 순전히 Jackson을 통해 동작하므로 단지 DruidModule이 반환하는 Jackson 모듈에 추가만 하면 돼요.
Complex Metric 추가
ComplexMetric 추가는 현재 버전에서 조금 지저분합니다. complex metric에 접근하는 방법은 ComplexMetrics.registerSerde() 메서드로 등록하는 것입니다. 특별한 Guice 처리 없이, 그냥 configure(Binder) 메서드에서 직렬화/역직렬화를 등록하면 됩니다.
새 Query 타입 추가
새 Query 타입 추가는 세 개의 인터페이스 구현이 필요해요.
org.apache.druid.query.Queryorg.apache.druid.query.QueryToolChestorg.apache.druid.query.QueryRunnerFactory
등록은 딥 스토리지 메커니즘과 같은 일반 전략을 사용합니다. 대략 다음과 같이 하면 되요.
DruidBinders.queryToolChestBinder(binder) .addBinding(SegmentMetadataQuery.class) .to(SegmentMetadataQueryQueryToolChest.class);DruidBinders.queryRunnerFactoryBinder(binder) .addBinding(SegmentMetadataQuery.class) .to(SegmentMetadataQueryRunnerFactory.class);
첫 번째는 SegmentMetadataQuery가 사용될 때 SegmentMetadataQueryQueryToolChest를 바인딩하고, 두 번째는 QueryRunnerFactory에 대해 같은 일을 합니다.
새 Jersey 리소스 추가
모듈에 새 Jersey 리소스를 추가하려면 모듈에서 리소스를 바인딩하기 위해 다음 코드를 호출해야 해요.
Jerseys.addResource(binder, NewResource.class);
새 Password Provider 구현 추가
org.apache.druid.metadata.PasswordProvider 인터페이스를 구현해야 합니다. Druid가 PasswordProvider를 사용하는 모든 곳에서 구현의 새 인스턴스가 생성되므로, 각 비밀번호를 가져오는 데 필요한 모든 정보가 객체 인스턴스화 시점에 제공되는지 확인하세요. org.apache.druid.initialization.DruidModule 구현에서 getJacksonModules는 대략 다음과 같아야 합니다.
return ImmutableList.of( new SimpleModule("SomePasswordProviderModule") .registerSubtypes( new NamedType(SomePasswordProvider.class, "some") ) );
여기서 SomePasswordProvider는 PasswordProvider 인터페이스의 구현이고, 예시로 org.apache.druid.metadata.EnvironmentVariablePasswordProvider를 볼 수 있어요.
새 DynamicConfigProvider 구현 추가
org.apache.druid.metadata.DynamicConfigProvider 인터페이스를 구현해야 합니다. Druid가 DynamicConfigProvider를 사용하는 모든 곳에서 구현의 새 인스턴스가 생성되므로, 모든 정보를 가져오는 데 필요한 정보가 객체 인스턴스화 시점에 제공되는지 확인하세요. org.apache.druid.initialization.DruidModule 구현에서 getJacksonModules는 대략 다음과 같아야 합니다.
return ImmutableList.of( new SimpleModule("SomeDynamicConfigProviderModule") .registerSubtypes( new NamedType(SomeDynamicConfigProvider.class, "some") ) );
여기서 SomeDynamicConfigProvider는 DynamicConfigProvider 인터페이스의 구현이고, 예시로 org.apache.druid.metadata.MapStringDynamicConfigProvider를 볼 수 있어요.
Transform 확장 추가
transform 확장을 만들려면 org.apache.druid.segment.transform.Transform 인터페이스를 구현하세요. org.apache.druid.segment.transform을 import하려면 druid-processing 패키지를 설치해야 합니다.
import com.fasterxml.jackson.annotation.JsonCreator;import com.fasterxml.jackson.annotation.JsonProperty;import org.apache.druid.segment.transform.RowFunction;import org.apache.druid.segment.transform.Transform;public class MyTransform implements Transform { private final String name; @JsonCreator public MyTransform( @JsonProperty("name") final String name ) { this.name = name; } @JsonProperty @Override public String getName() { return name; } @Override public RowFunction getRowFunction() { return new MyRowFunction(); } static class MyRowFunction implements RowFunction { @Override public Object eval(Row row) { return "transformed-value"; } }}
그런 다음 transform을 Jackson 모듈로 등록하세요.
import com.fasterxml.jackson.databind.Module;import com.fasterxml.jackson.databind.jsontype.NamedModule;import com.fasterxml.jackson.databind.module.SimpleModule;import com.google.inject.Binder;import com.google.common.collect.ImmutableList;import org.apache.druid.initialization.DruidModule;public class MyTransformModule implements DruidModule { @Override public List<? extends Module> getJacksonModules() { return return ImmutableList.of( new SimpleModule("MyTransformModule").registerSubtypes( new NamedType(MyTransform.class, "my-transform") ) ): } @Override public void configure(Binder binder) { }}
자신만의 커스텀 플러거블 Coordinator Duty 추가
코디네이터는 주기적으로 CoordinatorDuty라 불리는 작업들(새 세그먼트 로드, 세그먼트 밸런싱 등)을 실행해요. Druid 사용자는 Core Druid의 어떤 클래스도 수정하지 않고, Core Druid에 속하지 않는 커스텀 플러거블 코디네이터 duty를 추가할 수 있습니다. 사용자는 CoordinatorCustomDuty 인터페이스를 구현하고 JsonTypeName을 설정해서 자신만의 커스텀 코디네이터 duty를 작성하면 돼요. 다음으로 Module의 DruidModule#getJacksonModules()에서 커스텀 코디네이터를 서브타입으로 등록해야 합니다. 이 단계가 끝나면 다음 속성들로 커스텀 코디네이터 duty를 로드할 수 있습니다.
druid.coordinator.dutyGroups=[<GROUP_NAME_1>, <GROUP_NAME_2>, ...]druid.coordinator.<GROUP_NAME_1>.duties=[<DUTY_NAME_MATCHING_JSON_TYPE_NAME_1>, <DUTY_NAME_MATCHING_JSON_TYPE_NAME_2>, ...]druid.coordinator.<GROUP_NAME_1>.period=<GROUP_NAME_1_RUN_PERIOD>druid.coordinator.<GROUP_NAME_1>.duty.<DUTY_NAME_MATCHING_JSON_TYPE_NAME_1>.<SOME_CONFIG_1_KEY>=<SOME_CONFIG_1_VALUE>druid.coordinator.<GROUP_NAME_1>.duty.<DUTY_NAME_MATCHING_JSON_TYPE_NAME_1>.<SOME_CONFIG_2_KEY>=<SOME_CONFIG_2_VALUE>
새 플러거블 Coordinator duty 시스템에서, 코디네이터가 이미 하는 것과 유사하게 duty들을 그룹화할 수 있어요. duty들은 druid.coordinator.dutyGroups 목록의 요소에 따라 여러 그룹으로 묶입니다. 같은 그룹의 모든 duty는 druid.coordinator.<GROUP_NAME>.period로 설정된 동일한 실행 주기를 갖습니다. 현재 각 그룹에 대해 duty들을 순차적으로 실행하는 단일 스레드가 있습니다.
예를 들어 커스텀 코디네이터 duty 구현은 KillSupervisorsCustomDuty를, KillSupervisorsCustomDuty를 구성하는 데 사용할 수 있는 샘플 속성은 KillSupervisorsCustomDutyTest를 참고하세요.
확장의 HTTP 프록시 경유 데이터 라우팅
확장의 HttpClient가 HTTP 프록시를 통해 연결하도록 하는 기능을 추가할 수 있어요.
확장의 HTTP 클라이언트에 프록시 연결을 지원하려면:
- 확장의 HTTP config 클래스에
HttpClientProxyConfig를@JsonProperty로 추가하세요. - 확장의 모듈 클래스에서
HttpClientConfig에HttpProxyConfigconfig를 추가하세요. 예를 들어,config변수가 1단계의 확장 HTTP config일 때:
final HttpClientConfig.Builder builder = HttpClientConfig .builder() .withNumConnections(1) .withReadTimeout(config.getReadTimeout().toStandardDuration()) .withHttpProxyConfig(config.getProxyConfig());
확장을 다른 모든 Druid 확장과 함께 번들링
mvn install을 하면 Druid 확장은 distribution/target/ 아래에 있는 Druid tarball과 extensions 디렉토리 안에 패키징됩니다.
확장을 포함하고 싶다면 distribution/pom.xml에 확장의 maven 좌표를 인자로 추가할 수 있어요.
mvn install 동안 maven은 확장을 로컬 maven 저장소에 설치하고, 그런 다음 pull-deps를 호출해서 거기서 확장을 가져옵니다. 결국 distribution/target/extensions 아래와 Druid tarball 안에 확장이 보여야 합니다.
의존성 관리
공통으로 쓰이는 라이브러리를 끌어오는 확장에서는 라이브러리 충돌 관리가 만만치 않을 수 있어요. druid가 사용하는 버전과의 충돌을 막기 위해 provided 스코프로 지정할 것을 권장하는 라이브러리 group ID 목록입니다.
"org.apache.druid","com.metamx.druid","asm","org.ow2.asm","org.jboss.netty","com.google.guava","com.google.code.findbugs","com.google.protobuf","com.esotericsoftware.minlog","log4j","org.slf4j","commons-logging","org.eclipse.jetty","org.mortbay.jetty","com.sun.jersey","com.sun.jersey.contribs","common-beanutils","commons-codec","commons-lang","commons-cli","commons-io","javax.activation","org.apache.httpcomponents","org.apache.zookeeper","org.codehaus.jackson","com.fasterxml.jackson","com.fasterxml.jackson.core","com.fasterxml.jackson.dataformat","com.fasterxml.jackson.datatype","org.roaringbitmap","net.java.dev.jets3t"
더 자세한 내용은 org.apache.druid.cli.PullDependencies 문서를 참고하세요.
더 알아보기 (Learn more)
- Druid 개발 개요 — Druid 코드베이스 컴포넌트
- Druid 확장 코어 문서 — 내장 확장 목록
- 버전 관리 문서 — 확장 인터페이스의 버전 관리 정책