Tasks API
Tasks API
Apache Druid에서 태스크를 조회·제출·삭제하는 API 엔드포인트를 다룹니다. 태스크는 수집, 쿼리, 컴팩션 같은 작업을 완료하기 위해 Druid가 수행하는 개별 작업이에요.
출처: 문서
본문
이 문서는 Apache Druid의 태스크 조회, 제출, 삭제를 위한 API 엔드포인트를 설명합니다. 태스크(task)는 수집, 쿼리, 컴팩션 같은 작업을 완료하기 위해 Druid가 수행하는 개별 작업입니다.
이 문서에서 http://ROUTER_IP:ROUTER_PORT는 Router 서비스 주소와 포트를 위한 자리표시자입니다. 예를 들어 quickstart 구성에서는 http://localhost:8888을 사용하세요.
태스크 정보 및 조회
태스크 배열 가져오기
Druid 클러스터의 모든 태스크 배열을 가져옵니다. 각 태스크 객체는 ID, 상태, 연결된 데이터소스 및 기타 메타데이터에 대한 정보를 포함합니다. 응답 속성의 정의는 Tasks table을 참고하세요.
URL
GET /druid/indexer/v1/tasks
쿼리 파라미터
이 엔드포인트는 결과를 필터링하는 선택적 쿼리 파라미터 집합을 지원합니다.
| 파라미터 | 타입 | 설명 |
|---|---|---|
| state | String | 태스크 상태로 태스크 목록을 필터링합니다. 유효한 옵션은 running, complete, waiting, pending입니다. |
| datasource | String | Druid 데이터소스로 필터링된 태스크를 반환합니다. |
| createdTimeInterval | String (ISO-8601) | 지정된 간격 내에 생성된 태스크를 반환합니다. 간격 문자열 구분자로 _를 사용하세요. /는 사용하지 마세요. 예: 2023-06-27_2023-06-28. |
| max | Integer | 반환할 완료 태스크의 최대 개수. state가 complete로 설정된 경우에만 적용됩니다. |
| type | String | 태스크 타입으로 태스크를 필터링합니다. 자세한 내용은 task documentation을 참고하세요. |
응답
- 200 SUCCESS — 태스크 목록을 성공적으로 조회함
- 400 BAD REQUEST — 잘못된 state 쿼리 파라미터 값
- 500 SERVER ERROR — 잘못된 쿼리 파라미터
샘플 요청
다음 예시는 다음 쿼리 파라미터로 필터링된 태스크 목록을 가져오는 방법을 보여줍니다.
- State:
complete - Datasource:
wikipedia_api - 시간 간격: 2015-09-12에서 2015-09-13 사이
- 반환되는 최대 항목 수: 10
- 태스크 타입:
query_worker
cURL
curl "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/tasks/?state=complete&datasource=wikipedia_api&createdTimeInterval=2015-09-12_2015-09-13&max=10&type=query_worker"
HTTP
GET /druid/indexer/v1/tasks/?state=complete&datasource=wikipedia_api&createdTimeInterval=2015-09-12_2015-09-13&max=10&type=query_worker HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
[
{
"id": "query-223549f8-b993-4483-b028-1b0d54713cad-worker0_0",
"groupId": "query-223549f8-b993-4483-b028-1b0d54713cad",
"type": "query_worker",
"createdTime": "2023-06-22T22:11:37.012Z",
"queueInsertionTime": "1970-01-01T00:00:00.000Z",
"statusCode": "SUCCESS",
"status": "SUCCESS",
"runnerStatusCode": "NONE",
"duration": 17897,
"location": {
"host": "localhost",
"port": 8101,
"tlsPort": -1
},
"dataSource": "wikipedia_api",
"errorMsg": null
},
{
"id": "query-fa82fa40-4c8c-4777-b832-cabbee5f519f-worker0_0",
"groupId": "query-fa82fa40-4c8c-4777-b832-cabbee5f519f",
"type": "query_worker",
"createdTime": "2023-06-20T22:51:21.302Z",
"queueInsertionTime": "1970-01-01T00:00:00.000Z",
"statusCode": "SUCCESS",
"status": "SUCCESS",
"runnerStatusCode": "NONE",
"duration": 16911,
"location": {
"host": "localhost",
"port": 8101,
"tlsPort": -1
},
"dataSource": "wikipedia_api",
"errorMsg": null
},
{
"id": "query-5419da7a-b270-492f-90e6-920ecfba766a-worker0_0",
"groupId": "query-5419da7a-b270-492f-90e6-920ecfba766a",
"type": "query_worker",
"createdTime": "2023-06-20T22:45:53.909Z",
"queueInsertionTime": "1970-01-01T00:00:00.000Z",
"statusCode": "SUCCESS",
"status": "SUCCESS",
"runnerStatusCode": "NONE",
"duration": 17030,
"location": {
"host": "localhost",
"port": 8101,
"tlsPort": -1
},
"dataSource": "wikipedia_api",
"errorMsg": null
}
]
완료 태스크 배열 가져오기
Druid 클러스터의 완료된 태스크 배열을 가져옵니다. 이는 기능적으로 /druid/indexer/v1/tasks?state=complete와 동일합니다. 응답 속성의 정의는 Tasks table을 참고하세요.
URL
GET /druid/indexer/v1/completeTasks
쿼리 파라미터
이 엔드포인트는 결과를 필터링하는 선택적 쿼리 파라미터 집합을 지원합니다.
| 파라미터 | 타입 | 설명 |
|---|---|---|
| datasource | String | Druid 데이터소스로 필터링된 태스크를 반환합니다. |
| createdTimeInterval | String (ISO-8601) | 지정된 간격 내에 생성된 태스크를 반환합니다. 간격 문자열은 / 대신 _로 구분해야 합니다. 예: 2023-06-27_2023-06-28. |
| max | Integer | 반환할 완료 태스크의 최대 개수. state가 complete로 설정된 경우에만 적용됩니다. |
| type | String | 태스크 타입으로 태스크를 필터링합니다. 자세한 내용은 task documentation을 참고하세요. |
응답
- 200 SUCCESS — 완료 태스크 목록을 성공적으로 조회함
- 404 NOT FOUND — 잘못된 서비스로 요청이 전송됨
샘플 요청
cURL
curl "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/completeTasks"
HTTP
GET /druid/indexer/v1/completeTasks HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
[
{
"id": "query-223549f8-b993-4483-b028-1b0d54713cad-worker0_0",
"groupId": "query-223549f8-b993-4483-b028-1b0d54713cad",
"type": "query_worker",
"createdTime": "2023-06-22T22:11:37.012Z",
"queueInsertionTime": "1970-01-01T00:00:00.000Z",
"statusCode": "SUCCESS",
"status": "SUCCESS",
"runnerStatusCode": "NONE",
"duration": 17897,
"location": {
"host": "localhost",
"port": 8101,
"tlsPort": -1
},
"dataSource": "wikipedia_api",
"errorMsg": null
},
{
"id": "query-223549f8-b993-4483-b028-1b0d54713cad",
"groupId": "query-223549f8-b993-4483-b028-1b0d54713cad",
"type": "query_controller",
"createdTime": "2023-06-22T22:11:28.367Z",
"queueInsertionTime": "1970-01-01T00:00:00.000Z",
"statusCode": "SUCCESS",
"status": "SUCCESS",
"runnerStatusCode": "NONE",
"duration": 30317,
"location": {
"host": "localhost",
"port": 8100,
"tlsPort": -1
},
"dataSource": "wikipedia_api",
"errorMsg": null
}
]
실행 중인 태스크 배열 가져오기
Druid 클러스터의 실행 중인 태스크 객체 배열을 가져옵니다. 기능적으로 /druid/indexer/v1/tasks?state=running과 동일합니다. 응답 속성의 정의는 Tasks table을 참고하세요.
URL
GET /druid/indexer/v1/runningTasks
쿼리 파라미터
이 엔드포인트는 결과를 필터링하는 선택적 쿼리 파라미터 집합을 지원합니다.
| 파라미터 | 타입 | 설명 |
|---|---|---|
| datasource | String | Druid 데이터소스로 필터링된 태스크를 반환합니다. |
| createdTimeInterval | String (ISO-8601) | 지정된 간격 내에 생성된 태스크를 반환합니다. 간격 문자열은 / 대신 _로 구분해야 합니다. 예: 2023-06-27_2023-06-28. |
| max | Integer | 반환할 완료 태스크의 최대 개수. state가 complete로 설정된 경우에만 적용됩니다. |
| type | String | 태스크 타입으로 태스크를 필터링합니다. 자세한 내용은 task documentation을 참고하세요. |
응답
- 200 SUCCESS — 실행 중인 태스크 목록을 성공적으로 조회함
샘플 요청
cURL
curl "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/runningTasks"
HTTP
GET /druid/indexer/v1/runningTasks HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
[
{
"id": "query-32663269-ead9-405a-8eb6-0817a952ef47",
"groupId": "query-32663269-ead9-405a-8eb6-0817a952ef47",
"type": "query_controller",
"createdTime": "2023-06-22T22:54:43.170Z",
"queueInsertionTime": "2023-06-22T22:54:43.170Z",
"statusCode": "RUNNING",
"status": "RUNNING",
"runnerStatusCode": "RUNNING",
"duration": -1,
"location": {
"host": "localhost",
"port": 8100,
"tlsPort": -1
},
"dataSource": "wikipedia_api",
"errorMsg": null
}
]
대기 중인 태스크 배열 가져오기
Druid 클러스터의 대기 중인(waiting) 태스크 배열을 가져옵니다. 기능적으로 /druid/indexer/v1/tasks?state=waiting과 동일합니다. 응답 속성의 정의는 Tasks table을 참고하세요.
URL
GET /druid/indexer/v1/waitingTasks
쿼리 파라미터
이 엔드포인트는 결과를 필터링하는 선택적 쿼리 파라미터 집합을 지원합니다.
| 파라미터 | 타입 | 설명 |
|---|---|---|
| datasource | String | Druid 데이터소스로 필터링된 태스크를 반환합니다. |
| createdTimeInterval | String (ISO-8601) | 지정된 간격 내에 생성된 태스크를 반환합니다. 간격 문자열은 / 대신 _로 구분해야 합니다. 예: 2023-06-27_2023-06-28. |
| max | Integer | 반환할 완료 태스크의 최대 개수. state가 complete로 설정된 경우에만 적용됩니다. |
| type | String | 태스크 타입으로 태스크를 필터링합니다. 자세한 내용은 task documentation을 참고하세요. |
응답
- 200 SUCCESS — 대기 중인 태스크 목록을 성공적으로 조회함
샘플 요청
cURL
curl "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/waitingTasks"
HTTP
GET /druid/indexer/v1/waitingTasks HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
[
{
"id": "index_parallel_wikipedia_auto_biahcbmf_2023-06-26T21:08:05.216Z",
"groupId": "index_parallel_wikipedia_auto_biahcbmf_2023-06-26T21:08:05.216Z",
"type": "index_parallel",
"createdTime": "2023-06-26T21:08:05.217Z",
"queueInsertionTime": "1970-01-01T00:00:00.000Z",
"statusCode": "RUNNING",
"status": "RUNNING",
"runnerStatusCode": "WAITING",
"duration": -1,
"location": {
"host": null,
"port": -1,
"tlsPort": -1
},
"dataSource": "wikipedia_auto",
"errorMsg": null
},
{
"id": "index_parallel_wikipedia_auto_afggfiec_2023-06-26T21:08:05.546Z",
"groupId": "index_parallel_wikipedia_auto_afggfiec_2023-06-26T21:08:05.546Z",
"type": "index_parallel",
"createdTime": "2023-06-26T21:08:05.548Z",
"queueInsertionTime": "1970-01-01T00:00:00.000Z",
"statusCode": "RUNNING",
"status": "RUNNING",
"runnerStatusCode": "WAITING",
"duration": -1,
"location": {
"host": null,
"port": -1,
"tlsPort": -1
},
"dataSource": "wikipedia_auto",
"errorMsg": null
},
{
"id": "index_parallel_wikipedia_auto_jmmddihf_2023-06-26T21:08:06.644Z",
"groupId": "index_parallel_wikipedia_auto_jmmddihf_2023-06-26T21:08:06.644Z",
"type": "index_parallel",
"createdTime": "2023-06-26T21:08:06.671Z",
"queueInsertionTime": "1970-01-01T00:00:00.000Z",
"statusCode": "RUNNING",
"status": "RUNNING",
"runnerStatusCode": "WAITING",
"duration": -1,
"location": {
"host": null,
"port": -1,
"tlsPort": -1
},
"dataSource": "wikipedia_auto",
"errorMsg": null
}
]
pending 태스크 배열 가져오기
Druid 클러스터의 pending 태스크 배열을 가져옵니다. 기능적으로 /druid/indexer/v1/tasks?state=pending과 동일합니다. 응답 속성의 정의는 Tasks table을 참고하세요.
URL
GET /druid/indexer/v1/pendingTasks
쿼리 파라미터
이 엔드포인트는 결과를 필터링하는 선택적 쿼리 파라미터 집합을 지원합니다.
| 파라미터 | 타입 | 설명 |
|---|---|---|
| datasource | String | Druid 데이터소스로 필터링된 태스크를 반환합니다. |
| createdTimeInterval | String (ISO-8601) | 지정된 간격 내에 생성된 태스크를 반환합니다. 간격 문자열은 / 대신 _로 구분해야 합니다. 예: 2023-06-27_2023-06-28. |
| max | Integer | 반환할 완료 태스크의 최대 개수. state가 complete로 설정된 경우에만 적용됩니다. |
| type | String | 태스크 타입으로 태스크를 필터링합니다. 자세한 내용은 task documentation을 참고하세요. |
응답
- 200 SUCCESS — pending 태스크 목록을 성공적으로 조회함
샘플 요청
cURL
curl "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/pendingTasks"
HTTP
GET /druid/indexer/v1/pendingTasks HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
[
{
"id": "query-7b37c315-50a0-4b68-aaa8-b1ef1f060e67",
"groupId": "query-7b37c315-50a0-4b68-aaa8-b1ef1f060e67",
"type": "query_controller",
"createdTime": "2023-06-23T19:53:06.037Z",
"queueInsertionTime": "2023-06-23T19:53:06.037Z",
"statusCode": "RUNNING",
"status": "RUNNING",
"runnerStatusCode": "PENDING",
"duration": -1,
"location": {
"host": null,
"port": -1,
"tlsPort": -1
},
"dataSource": "wikipedia_api",
"errorMsg": null
},
{
"id": "query-544f0c41-f81d-4504-b98b-f9ab8b36ef36",
"groupId": "query-544f0c41-f81d-4504-b98b-f9ab8b36ef36",
"type": "query_controller",
"createdTime": "2023-06-23T19:53:06.616Z",
"queueInsertionTime": "2023-06-23T19:53:06.616Z",
"statusCode": "RUNNING",
"status": "RUNNING",
"runnerStatusCode": "PENDING",
"duration": -1,
"location": {
"host": null,
"port": -1,
"tlsPort": -1
},
"dataSource": "wikipedia_api",
"errorMsg": null
}
]
태스크 페이로드 가져오기
태스크 ID가 주어지면 태스크의 페이로드를 가져옵니다. 태스크 ID와 태스크 실행과 연결된 구성 세부 정보 및 관련 스펙을 포함하는 페이로드가 담긴 JSON 객체를 반환합니다.
URL
GET /druid/indexer/v1/task/{taskId}
응답
- 200 SUCCESS — 태스크 페이로드를 성공적으로 조회함
- 404 NOT FOUND — ID로 태스크를 찾을 수 없음
샘플 요청
다음 예시는 지정된 ID index_parallel_wikipedia_short_iajoonnd_2023-07-07T17:53:12.174Z의 태스크 페이로드를 가져오는 방법을 보여줍니다.
cURL
curl "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/task/index_parallel_wikipedia_short_iajoonnd_2023-07-07T17:53:12.174Z"
HTTP
GET /druid/indexer/v1/task/index_parallel_wikipedia_short_iajoonnd_2023-07-07T17:53:12.174Z HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
{
"task": "index_parallel_wikipedia_short_iajoonnd_2023-07-07T17:53:12.174Z",
"payload": {
"type": "index_parallel",
"id": "index_parallel_wikipedia_short_iajoonnd_2023-07-07T17:53:12.174Z",
"groupId": "index_parallel_wikipedia_short_iajoonnd_2023-07-07T17:53:12.174Z",
"resource": {
"availabilityGroup": "index_parallel_wikipedia_short_iajoonnd_2023-07-07T17:53:12.174Z",
"requiredCapacity": 1
},
"spec": {
"dataSchema": {
"dataSource": "wikipedia_short",
"timestampSpec": {
"column": "time",
"format": "iso",
"missingValue": null
},
"dimensionsSpec": {
"dimensions": [
{
"type": "string",
"name": "cityName",
"multiValueHandling": "SORTED_ARRAY",
"createBitmapIndex": true
},
{
"type": "string",
"name": "countryName",
"multiValueHandling": "SORTED_ARRAY",
"createBitmapIndex": true
},
{
"type": "string",
"name": "regionName",
"multiValueHandling": "SORTED_ARRAY",
"createBitmapIndex": true
}
],
"dimensionExclusions": [
"__time",
"time"
],
"includeAllDimensions": false,
"useSchemaDiscovery": false
},
"metricsSpec": [],
"granularitySpec": {
"type": "uniform",
"segmentGranularity": "DAY",
"queryGranularity": {
"type": "none"
},
"rollup": false,
"intervals": [
"2015-09-12T00:00:00.000Z/2015-09-13T00:00:00.000Z"
]
},
"transformSpec": {
"filter": null,
"transforms": []
}
},
"ioConfig": {
"type": "index_parallel",
"inputSource": {
"type": "local",
"baseDir": "quickstart/tutorial",
"filter": "wikiticker-2015-09-12-sampled.json.gz"
},
"inputFormat": {
"type": "json"
},
"appendToExisting": false,
"dropExisting": false
},
"tuningConfig": {
"type": "index_parallel",
"maxRowsPerSegment": 5000000,
"appendableIndexSpec": {
"type": "onheap",
"preserveExistingMetrics": false
},
"maxRowsInMemory": 25000,
"maxBytesInMemory": 0,
"skipBytesInMemoryOverheadCheck": false,
"maxTotalRows": null,
"numShards": null,
"splitHintSpec": null,
"partitionsSpec": {
"type": "dynamic",
"maxRowsPerSegment": 5000000,
"maxTotalRows": null
},
"indexSpec": {
"bitmap": {
"type": "roaring"
},
"dimensionCompression": "lz4",
"stringDictionaryEncoding": {
"type": "utf8"
},
"metricCompression": "lz4",
"longEncoding": "longs"
},
"indexSpecForIntermediatePersists": {
"bitmap": {
"type": "roaring"
},
"dimensionCompression": "lz4",
"stringDictionaryEncoding": {
"type": "utf8"
},
"metricCompression": "lz4",
"longEncoding": "longs"
},
"maxPendingPersists": 0,
"forceGuaranteedRollup": false,
"reportParseExceptions": false,
"pushTimeout": 0,
"segmentWriteOutMediumFactory": null,
"maxNumConcurrentSubTasks": 1,
"maxRetry": 3,
"taskStatusCheckPeriodMs": 1000,
"chatHandlerTimeout": "PT10S",
"chatHandlerNumRetries": 5,
"maxNumSegmentsToMerge": 100,
"totalNumMergeTasks": 10,
"logParseExceptions": false,
"maxParseExceptions": 2147483647,
"maxSavedParseExceptions": 0,
"maxColumnsToMerge": -1,
"awaitSegmentAvailabilityTimeoutMillis": 0,
"maxAllowedLockCount": -1,
"partitionDimensions": []
}
},
"context": {
"forceTimeChunkLock": true,
"useLineageBasedSegmentAllocation": true
},
"dataSource": "wikipedia_short"
}
}
태스크 상태 가져오기
태스크 ID가 주어지면 태스크의 상태를 가져옵니다. 태스크의 상태 코드, runner 상태, 태스크 타입, 데이터소스 및 기타 관련 메타데이터가 담긴 JSON 객체를 반환합니다.
URL
GET /druid/indexer/v1/task/{taskId}/status
응답
- 200 SUCCESS — 태스크 상태를 성공적으로 조회함
- 404 NOT FOUND — ID로 태스크를 찾을 수 없음
샘플 요청
다음 예시는 지정된 ID query-223549f8-b993-4483-b028-1b0d54713cad의 태스크 상태를 가져오는 방법을 보여줍니다.
cURL
curl "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/task/query-223549f8-b993-4483-b028-1b0d54713cad/status"
HTTP
GET /druid/indexer/v1/task/query-223549f8-b993-4483-b028-1b0d54713cad/status HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
{
"task": "query-223549f8-b993-4483-b028-1b0d54713cad",
"status": {
"id": "query-223549f8-b993-4483-b028-1b0d54713cad",
"groupId": "query-223549f8-b993-4483-b028-1b0d54713cad",
"type": "query_controller",
"createdTime": "2023-06-22T22:11:28.367Z",
"queueInsertionTime": "1970-01-01T00:00:00.000Z",
"statusCode": "RUNNING",
"status": "RUNNING",
"runnerStatusCode": "RUNNING",
"duration": -1,
"location": {"host": "localhost", "port": 8100, "tlsPort": -1},
"dataSource": "wikipedia_api",
"errorMsg": null
}
}
태스크 세그먼트 가져오기
참고: 이 API는 더 이상 지원되지 않으며 항상 404 응답을 반환합니다. 태스크가 커밋한 세그먼트 ID를 식별하려면 대신 메트릭
segment/added/bytes를 사용하세요.
URL
GET /druid/indexer/v1/task/{taskId}/segments
응답
- 404 NOT FOUND
{
"error": "Segment IDs committed by a task action are not persisted anymore. Use the metric 'segment/added/bytes' to identify the segments created by a task."
}
샘플 요청
다음 예시는 지정된 ID query-52a8aafe-7265-4427-89fe-dc51275cc470의 태스크 세그먼트를 가져오는 방법을 보여줍니다.
cURL
curl "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/task/query-52a8aafe-7265-4427-89fe-dc51275cc470/reports"
HTTP
GET /druid/indexer/v1/task/query-52a8aafe-7265-4427-89fe-dc51275cc470/reports HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
성공한 요청은 200 OK 응답과 태스크 세그먼트 배열을 반환합니다.
태스크 로그 가져오기
태스크와 연결된 이벤트 로그를 가져옵니다. 태스크의 수명주기 동안 기록된 로그 이벤트 목록을 반환합니다. 이 엔드포인트는 오류나 경고를 포함한 태스크 실행에 대한 정보를 제공하는 데 유용합니다.
태스크 로그는 Middle Manager/Indexer 또는 장기 저장소에서 자동으로 가져옵니다. 참고로 Task logs를 참고하세요.
URL
GET /druid/indexer/v1/task/{taskId}/log
쿼리 파라미터
- offset (선택) — 타입: Int. 응답에서 처음 전달된 개수의 항목을 제외합니다.
응답
- 200 SUCCESS — 태스크 로그를 성공적으로 조회함
샘플 요청
다음 예시는 지정된 ID index_kafka_social_media_0e905aa31037879_nommnaeg의 태스크 로그를 가져오는 방법을 보여줍니다.
cURL
curl "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/task/index_kafka_social_media_0e905aa31037879_nommnaeg/log"
HTTP
GET /druid/indexer/v1/task/index_kafka_social_media_0e905aa31037879_nommnaeg/log HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
2023-07-03T22:11:17,891 INFO [qtp1251996697-122] org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner - Sequence[index_kafka_social_media_0e905aa31037879_0] end offsets updated from [{0=9223372036854775807}] to [{0=230985}].
2023-07-03T22:11:17,900 INFO [qtp1251996697-122] org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner - Saved sequence metadata to disk: [SequenceMetadata{sequenceId=0, sequenceName='index_kafka_social_media_0e905aa31037879_0', assignments=[0], startOffsets={0=230985}, exclusiveStartPartitions=[], endOffsets={0=230985}, sentinel=false, checkpointed=true}]
2023-07-03T22:11:17,901 INFO [task-runner-0-priority-0] org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner - Received resume command, resuming ingestion.
2023-07-03T22:11:17,901 INFO [task-runner-0-priority-0] org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner - Finished reading partition[0], up to[230985].
2023-07-03T22:11:17,902 INFO [task-runner-0-priority-0] org.apache.kafka.clients.consumer.internals.ConsumerCoordinator - [Consumer clientId=consumer-kafka-supervisor-dcanhmig-1, groupId=kafka-supervisor-dcanhmig] Resetting generation and member id due to: consumer pro-actively leaving the group
2023-07-03T22:11:17,902 INFO [task-runner-0-priority-0] org.apache.kafka.clients.consumer.internals.ConsumerCoordinator - [Consumer clientId=consumer-kafka-supervisor-dcanhmig-1, groupId=kafka-supervisor-dcanhmig] Request joining group due to: consumer pro-actively leaving the group
2023-07-03T22:11:17,902 INFO [task-runner-0-priority-0] org.apache.kafka.clients.consumer.KafkaConsumer - [Consumer clientId=consumer-kafka-supervisor-dcanhmig-1, groupId=kafka-supervisor-dcanhmig] Unsubscribed all topics or patterns and assigned partitions
2023-07-03T22:11:17,912 INFO [task-runner-0-priority-0] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Persisted rows[0] and (estimated) bytes[0]
2023-07-03T22:11:17,916 INFO [[index_kafka_social_media_0e905aa31037879_nommnaeg]-appenderator-persist] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Flushed in-memory data with commit metadata [AppenderatorDriverMetadata{segments={}, lastSegmentIds={}, callerMetadata={nextPartitions=SeekableStreamEndSequenceNumbers{stream='social_media', partitionSequenceNumberMap={0=230985}}}}] for segments:
2023-07-03T22:11:17,917 INFO [[index_kafka_social_media_0e905aa31037879_nommnaeg]-appenderator-persist] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Persisted stats: processed rows: [0], persisted rows[0], sinks: [0], total fireHydrants (across sinks): [0], persisted fireHydrants (across sinks): [0]
2023-07-03T22:11:17,919 INFO [task-runner-0-priority-0] org.apache.druid.segment.realtime.appenderator.BaseAppenderatorDriver - Pushing [0] segments in background
2023-07-03T22:11:17,921 INFO [task-runner-0-priority-0] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Persisted rows[0] and (estimated) bytes[0]
2023-07-03T22:11:17,924 INFO [[index_kafka_social_media_0e905aa31037879_nommnaeg]-appenderator-persist] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Flushed in-memory data with commit metadata [AppenderatorDriverMetadata{segments={}, lastSegmentIds={}, callerMetadata={nextPartitions=SeekableStreamStartSequenceNumbers{stream='social_media', partitionSequenceNumberMap={0=230985}, exclusivePartitions=[]}, publishPartitions=SeekableStreamEndSequenceNumbers{stream='social_media', partitionSequenceNumberMap={0=230985}}}}] for segments:
2023-07-03T22:11:17,924 INFO [[index_kafka_social_media_0e905aa31037879_nommnaeg]-appenderator-persist] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Persisted stats: processed rows: [0], persisted rows[0], sinks: [0], total fireHydrants (across sinks): [0], persisted fireHydrants (across sinks): [0]
2023-07-03T22:11:17,925 INFO [[index_kafka_social_media_0e905aa31037879_nommnaeg]-appenderator-merge] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Preparing to push (stats): processed rows: [0], sinks: [0], fireHydrants (across sinks): [0]
2023-07-03T22:11:17,925 INFO [[index_kafka_social_media_0e905aa31037879_nommnaeg]-appenderator-merge] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Push complete...
2023-07-03T22:11:17,929 INFO [[index_kafka_social_media_0e905aa31037879_nommnaeg]-publish] org.apache.druid.indexing.seekablestream.SequenceMetadata - With empty segment set, start offsets [SeekableStreamStartSequenceNumbers{stream='social_media', partitionSequenceNumberMap={0=230985}, exclusivePartitions=[]}] and end offsets [SeekableStreamEndSequenceNumbers{stream='social_media', partitionSequenceNumberMap={0=230985}}] are the same, skipping metadata commit.
2023-07-03T22:11:17,930 INFO [[index_kafka_social_media_0e905aa31037879_nommnaeg]-publish] org.apache.druid.segment.realtime.appenderator.BaseAppenderatorDriver - Published [0] segments with commit metadata [{nextPartitions=SeekableStreamStartSequenceNumbers{stream='social_media', partitionSequenceNumberMap={0=230985}, exclusivePartitions=[]}, publishPartitions=SeekableStreamEndSequenceNumbers{stream='social_media', partitionSequenceNumberMap={0=230985}}}]
2023-07-03T22:11:17,930 INFO [[index_kafka_social_media_0e905aa31037879_nommnaeg]-publish] org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner - Published 0 segments for sequence [index_kafka_social_media_0e905aa31037879_0] with metadata [AppenderatorDriverMetadata{segments={}, lastSegmentIds={}, callerMetadata={nextPartitions=SeekableStreamStartSequenceNumbers{stream='social_media', partitionSequenceNumberMap={0=230985}, exclusivePartitions=[]}, publishPartitions=SeekableStreamEndSequenceNumbers{stream='social_media', partitionSequenceNumberMap={0=230985}}}}].
2023-07-03T22:11:17,931 INFO [[index_kafka_social_media_0e905aa31037879_nommnaeg]-publish] org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner - Saved sequence metadata to disk: []
2023-07-03T22:11:17,932 INFO [task-runner-0-priority-0] org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner - Handoff complete for segments:
2023-07-03T22:11:17,932 INFO [task-runner-0-priority-0] org.apache.kafka.clients.consumer.internals.ConsumerCoordinator - [Consumer clientId=consumer-kafka-supervisor-dcanhmig-1, groupId=kafka-supervisor-dcanhmig] Resetting generation and member id due to: consumer pro-actively leaving the group
2023-07-03T22:11:17,932 INFO [task-runner-0-priority-0] org.apache.kafka.clients.consumer.internals.ConsumerCoordinator - [Consumer clientId=consumer-kafka-supervisor-dcanhmig-1, groupId=kafka-supervisor-dcanhmig] Request joining group due to: consumer pro-actively leaving the group
2023-07-03T22:11:17,933 INFO [task-runner-0-priority-0] org.apache.kafka.common.metrics.Metrics - Metrics scheduler closed
2023-07-03T22:11:17,933 INFO [task-runner-0-priority-0] org.apache.kafka.common.metrics.Metrics - Closing reporter org.apache.kafka.common.metrics.JmxReporter
2023-07-03T22:11:17,933 INFO [task-runner-0-priority-0] org.apache.kafka.common.metrics.Metrics - Metrics reporters closed
2023-07-03T22:11:17,935 INFO [task-runner-0-priority-0] org.apache.kafka.common.utils.AppInfoParser - App info kafka.consumer for consumer-kafka-supervisor-dcanhmig-1 unregistered
2023-07-03T22:11:17,936 INFO [task-runner-0-priority-0] org.apache.druid.curator.announcement.PathChildrenAnnouncer - Unannouncing [/druid/internal-discovery/PEON/localhost:8100]
2023-07-03T22:11:17,972 INFO [task-runner-0-priority-0] org.apache.druid.curator.discovery.CuratorDruidNodeAnnouncer - Unannounced self [{"druidNode":{"service":"druid/middleManager","host":"localhost","bindOnHost":false,"plaintextPort":8100,"port":-1,"tlsPort":-1,"enablePlaintextPort":true,"enableTlsPort":false},"nodeType":"peon","services":{"dataNodeService":{"type":"dataNodeService","tier":"_default_tier","maxSize":0,"type":"indexer-executor","serverType":"indexer-executor","priority":0},"lookupNodeService":{"type":"lookupNodeService","lookupTier":"__default"}}}].
2023-07-03T22:11:17,972 INFO [task-runner-0-priority-0] org.apache.druid.curator.announcement.PathChildrenAnnouncer - Unannouncing [/druid/announcements/localhost:8100]
2023-07-03T22:11:17,996 INFO [task-runner-0-priority-0] org.apache.druid.indexing.worker.executor.ExecutorLifecycle - Task completed with status: {
"id" : "index_kafka_social_media_0e905aa31037879_nommnaeg",
"status" : "SUCCESS",
"duration" : 3601130,
"errorMsg" : null,
"location" : {
"host" : null,
"port" : -1,
"tlsPort" : -1
}
}
2023-07-03T22:11:17,998 INFO [main] org.apache.druid.java.util.common.lifecycle.Lifecycle - Stopping lifecycle [module] stage [ANNOUNCEMENTS]
2023-07-03T22:11:18,005 INFO [main] org.apache.druid.java.util.common.lifecycle.Lifecycle - Stopping lifecycle [module] stage [SERVER]
2023-07-03T22:11:18,009 INFO [main] org.eclipse.jetty.server.AbstractConnector - Stopped ServerConnector@6491006{HTTP/1.1, (http/1.1)}{0.0.0.0:8100}
2023-07-03T22:11:18,009 INFO [main] org.eclipse.jetty.server.session - node0 Stopped scavenging
2023-07-03T22:11:18,012 INFO [main] org.eclipse.jetty.server.handler.ContextHandler - Stopped o.e.j.s.ServletContextHandler@742aa00a{/,null,STOPPED}
2023-07-03T22:11:18,014 INFO [main] org.apache.druid.java.util.common.lifecycle.Lifecycle - Stopping lifecycle [module] stage [NORMAL]
2023-07-03T22:11:18,014 INFO [main] org.apache.druid.server.coordination.ZkCoordinator - Stopping ZkCoordinator for [DruidServerMetadata{name='localhost:8100', hostAndPort='localhost:8100', hostAndTlsPort='null', maxSize=0, tier='_default_tier', type=indexer-executor, priority=0}]
2023-07-03T22:11:18,014 INFO [main] org.apache.druid.server.coordination.SegmentLoadDropHandler - Stopping...
2023-07-03T22:11:18,014 INFO [main] org.apache.druid.server.coordination.SegmentLoadDropHandler - Stopped.
2023-07-03T22:11:18,014 INFO [main] org.apache.druid.indexing.overlord.SingleTaskBackgroundRunner - Starting graceful shutdown of task[index_kafka_social_media_0e905aa31037879_nommnaeg].
2023-07-03T22:11:18,014 INFO [main] org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner - Stopping forcefully (status: [PUBLISHING])
2023-07-03T22:11:18,019 INFO [LookupExtractorFactoryContainerProvider-MainThread] org.apache.druid.query.lookup.LookupReferencesManager - Lookup Management loop exited. Lookup notices are not handled anymore.
2023-07-03T22:11:18,020 INFO [main] org.apache.druid.query.lookup.LookupReferencesManager - Closed lookup [name].
2023-07-03T22:11:18,020 INFO [Curator-Framework-0] org.apache.curator.framework.imps.CuratorFrameworkImpl - backgroundOperationsLoop exiting
2023-07-03T22:11:18,147 INFO [main] org.apache.zookeeper.ZooKeeper - Session: 0x1000097ceaf0007 closed
2023-07-03T22:11:18,147 INFO [main-EventThread] org.apache.zookeeper.ClientCnxn - EventThread shut down for session: 0x1000097ceaf0007
2023-07-03T22:11:18,151 INFO [main] org.apache.druid.java.util.common.lifecycle.Lifecycle - Stopping lifecycle [module] stage [INIT]
Finished peon task
태스크 완료 리포트 가져오기
태스크의 완료 리포트를 가져옵니다. 수집된 행 수와 Druid가 발생시킨 parse exception에 대한 정보가 담긴 JSON 객체를 반환합니다.
URL
GET /druid/indexer/v1/task/{taskId}/reports
응답
- 200 SUCCESS — 태스크 리포트를 성공적으로 조회함
샘플 요청
다음 예시는 지정된 ID query-52a8aafe-7265-4427-89fe-dc51275cc470의 태스크 완료 리포트를 가져오는 방법을 보여줍니다.
cURL
curl "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/task/query-52a8aafe-7265-4427-89fe-dc51275cc470/reports"
HTTP
GET /druid/indexer/v1/task/query-52a8aafe-7265-4427-89fe-dc51275cc470/reports HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
{
"ingestionStatsAndErrors": {
"type": "ingestionStatsAndErrors",
"taskId": "query-52a8aafe-7265-4427-89fe-dc51275cc470",
"payload": {
"ingestionState": "COMPLETED",
"unparseableEvents": {},
"rowStats": {
"determinePartitions": {
"processed": 0,
"processedBytes": 0,
"processedWithError": 0,
"thrownAway": 0,
"unparseable": 0
},
"buildSegments": {
"processed": 39244,
"processedBytes": 17106256,
"processedWithError": 0,
"thrownAway": 0,
"unparseable": 0
}
},
"errorMsg": null,
"segmentAvailabilityConfirmed": false,
"segmentAvailabilityWaitTimeMs": 0
}
}
}
태스크 작업 (Task operations)
태스크 제출
JSON 기반 수집 스펙이나 supervisor 스펙을 Overlord에 제출합니다. 제출된 태스크의 태스크 ID를 반환합니다. 수집 스펙 생성에 대한 자세한 내용은 ingestion spec reference를 참고하세요.
대부분의 배치 수집 유스케이스에서는 JSON 기반 배치 수집 대신 SQL-ingestion API를 사용해야 합니다.
URL
POST /druid/indexer/v1/task
응답
- 200 SUCCESS — 태스크를 성공적으로 제출함
- 400 BAD REQUEST — 쿼리에 정보 누락
- 415 UNSUPPORTED MEDIA TYPE — 잘못된 요청 본문 미디어 타입
- 500 Server Error — 요청 본문의 예상치 못한 토큰 또는 문자
샘플 요청
다음 요청은 "wikipedia auto"라는 이름의 데이터소스를 생성하는 태스크를 제출하는 예시입니다.
cURL
curl "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/task" \
--header 'Content-Type: application/json' \
--data '{
"type" : "index_parallel",
"spec" : {
"dataSchema" : {
"dataSource" : "wikipedia_auto",
"timestampSpec": {
"column": "time",
"format": "iso"
},
"dimensionsSpec" : {
"useSchemaDiscovery": true
},
"metricsSpec" : [],
"granularitySpec" : {
"type" : "uniform",
"segmentGranularity" : "day",
"queryGranularity" : "none",
"intervals" : ["2015-09-12/2015-09-13"],
"rollup" : false
}
},
"ioConfig" : {
"type" : "index_parallel",
"inputSource" : {
"type" : "local",
"baseDir" : "quickstart/tutorial/",
"filter" : "wikiticker-2015-09-12-sampled.json.gz"
},
"inputFormat" : {
"type" : "json"
},
"appendToExisting" : false
},
"tuningConfig" : {
"type" : "index_parallel",
"maxRowsPerSegment" : 5000000,
"maxRowsInMemory" : 25000
}
}
}'
HTTP
POST /druid/indexer/v1/task HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
Content-Type: application/json
Content-Length: 952
{
"type" : "index_parallel",
"spec" : {
"dataSchema" : {
"dataSource" : "wikipedia_auto",
"timestampSpec": {
"column": "time",
"format": "iso"
},
"dimensionsSpec" : {
"useSchemaDiscovery": true
},
"metricsSpec" : [],
"granularitySpec" : {
"type" : "uniform",
"segmentGranularity" : "day",
"queryGranularity" : "none",
"intervals" : ["2015-09-12/2015-09-13"],
"rollup" : false
}
},
"ioConfig" : {
"type" : "index_parallel",
"inputSource" : {
"type" : "local",
"baseDir" : "quickstart/tutorial/",
"filter" : "wikiticker-2015-09-12-sampled.json.gz"
},
"inputFormat" : {
"type" : "json"
},
"appendToExisting" : false
},
"tuningConfig" : {
"type" : "index_parallel",
"maxRowsPerSegment" : 5000000,
"maxRowsInMemory" : 25000
}
}
}
샘플 응답
{
"task": "index_parallel_wikipedia_odofhkle_2023-06-23T21:07:28.226Z"
}
태스크 종료
아직 완료되지 않았다면 태스크를 종료합니다. 성공적으로 종료된 태스크의 ID가 담긴 JSON 객체를 반환합니다.
URL
POST /druid/indexer/v1/task/{taskId}/shutdown
응답
- 200 SUCCESS — 태스크를 성공적으로 종료함
- 404 NOT FOUND — ID로 태스크를 찾을 수 없거나 태스크가 더 이상 실행 중이 아님
샘플 요청
다음 요청은 ID가 query-52as 8aafe-7265-4427-89fe-dc51275cc470인 태스크를 종료하는 방법을 보여줍니다.
cURL
curl --request POST "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/task/query-52as 8aafe-7265-4427-89fe-dc51275cc470/shutdown"
HTTP
POST /druid/indexer/v1/task/query-52as 8aafe-7265-4427-89fe-dc51275cc470/shutdown HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
{
"task": "query-577a83dd-a14e-4380-bd01-c942b781236b"
}
데이터소스의 모든 태스크 종료
지정된 데이터소스의 모든 태스크를 종료합니다. 성공하면 태스크가 종료된 데이터소스의 이름이 담긴 JSON 객체를 반환합니다.
URL
POST /druid/indexer/v1/datasources/{datasource}/shutdownAllTasks
응답
- 200 SUCCESS — 태스크를 성공적으로 종료함
- 404 NOT FOUND — 오류 또는 데이터소스에 실행 중인 태스크가 없음
샘플 요청
다음 요청은 wikipedia_auto 데이터소스의 모든 태스크를 종료하는 예시입니다.
cURL
curl --request POST "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/datasources/wikipedia_auto/shutdownAllTasks"
HTTP
POST /druid/indexer/v1/datasources/wikipedia_auto/shutdownAllTasks HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
{
"dataSource": "wikipedia_api"
}
태스크 관리 (Task management)
태스크의 상태 객체 가져오기
요청 본문의 태스크 ID 문자열 목록에 대한 태스크 상태 객체 목록을 가져옵니다. 각 태스크의 상태, 기간, 위치 및 오류 메시지가 담긴 JSON 객체 집합을 반환합니다.
URL
POST /druid/indexer/v1/taskStatus
응답
- 200 SUCCESS — 상태 객체를 성공적으로 조회함
- 415 UNSUPPORTED MEDIA TYPE — 요청 본문이 없거나 요청 본문 타입이 잘못됨
샘플 요청
다음 요청은 태스크 ID index_parallel_wikipedia_auto_jndhkpbo_2023-06-26T17:23:05.308Z와 index_parallel_wikipedia_auto_jbgiianh_2023-06-26T23:17:56.769Z의 상태 객체를 가져오는 예시입니다.
cURL
curl "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/taskStatus" \
--header 'Content-Type: application/json' \
--data '["index_parallel_wikipedia_auto_jndhkpbo_2023-06-26T17:23:05.308Z","index_parallel_wikipedia_auto_jbgiianh_2023-06-26T23:17:56.769Z"]'
HTTP
POST /druid/indexer/v1/taskStatus HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
Content-Type: application/json
Content-Length: 134
["index_parallel_wikipedia_auto_jndhkpbo_2023-06-26T17:23:05.308Z", "index_parallel_wikipedia_auto_jbgiianh_2023-06-26T23:17:56.769Z"]
샘플 응답
{
"index_parallel_wikipedia_auto_jbgiianh_2023-06-26T23:17:56.769Z": {
"id": "index_parallel_wikipedia_auto_jbgiianh_2023-06-26T23:17:56.769Z",
"status": "SUCCESS",
"duration": 10630,
"errorMsg": null,
"location": {
"host": "localhost",
"port": 8100,
"tlsPort": -1
}
},
"index_parallel_wikipedia_auto_jndhkpbo_2023-06-26T17:23:05.308Z": {
"id": "index_parallel_wikipedia_auto_jndhkpbo_2023-06-26T17:23:05.308Z",
"status": "SUCCESS",
"duration": 11012,
"errorMsg": null,
"location": {
"host": "localhost",
"port": 8100,
"tlsPort": -1
}
}
}
데이터소스의 pending 세그먼트 정리
데이터소스에 대해 메타데이터 스토리지의 pending segments 테이블을 수동으로 정리합니다. pending segments 테이블에서 삭제된 행 수에 대한 numDeleted가 담긴 JSON 객체 응답을 반환합니다. 이 API는 이 작업을 주기적으로 자동 수행하는 druid.coordinator.kill.pendingSegments.on Coordinator 설정이 사용합니다.
URL
DELETE /druid/indexer/v1/pendingSegments/{datasource}
응답
- 200 SUCCESS — pending 세그먼트를 성공적으로 삭제함
샘플 요청
다음 요청은 wikipedia_api 데이터소스의 pending 세그먼트를 정리하는 예시입니다.
cURL
curl --request DELETE "http://ROUTER_IP:ROUTER_PORT/druid/indexer/v1/pendingSegments/wikipedia_api"
HTTP
DELETE /druid/indexer/v1/pendingSegments/wikipedia_api HTTP/1.1
Host: http://ROUTER_IP:ROUTER_PORT
샘플 응답
{
"numDeleted": 2
}
더 알아보기 (Learn more)
- 태스크 유형과 수집 스펙에 대한 자세한 내용은 ingestion spec 문서를 참고하세요.
- supervisor 관리에 대한 내용은 supervisor API 문서를 확인해 보세요.