Ranger 의 감사 기록은 배포판에 따라 Solr 컬렉션(ranger_audits) · HDFS · 외부 저장소에 쌓인다. Solr 에 있는 기록을 Kafka 로 보내 다른 시스템에서 소비하려는 경우, Solr 에는 Kafka 로 내보내는 기능이 없으므로 중간에 읽어서 보내는 주체가 필요하다.
Kafka Connect 에는 Solr 싱크 커넥터는 있으나 공식 소스 커넥터가 없다. 결국 선택지는 주기적으로 질의해 보내는 방식이다.
컬렉션 목록과 상태를 먼저 확인한다.
SOLR=http://solr01.example.com:8983
curl -s "${SOLR}/solr/admin/collections?action=LIST&wt=json" | jq -r '.collections[]'
curl -s "${SOLR}/solr/admin/collections?action=CLUSTERSTATUS&wt=json" | jq '.cluster.collections | keys'
기본 질의 형태는 다음과 같다.
curl -s "${SOLR}/solr/ranger_audits/select" \
--data-urlencode 'q=*:*' \
--data-urlencode 'fq=evtTime:[NOW-1DAY TO NOW]' \
--data-urlencode 'sort=evtTime asc' \
--data-urlencode 'rows=1000' \
--data-urlencode 'wt=json' | jq '.response.numFound'
| 파라미터 | 의미 |
|---|---|
q |
질의. 전체는 *:* |
fq |
필터 질의. 캐시되므로 범위 조건은 여기에 둔다 |
fl |
돌려받을 필드 목록 |
sort |
정렬 |
rows · start |
개수와 시작 위치 |
wt |
응답 형식. json |
건수가 많으면 start 로 페이지를 넘기면 안 된다. 깊은 페이징은 급격히 느려진다. cursorMark 를 쓴다. 이때 sort 에 유일 키를 포함해야 한다.
curl -s "${SOLR}/solr/ranger_audits/select" \
--data-urlencode 'q=*:*' \
--data-urlencode 'fq=evtTime:[NOW-1HOUR TO NOW]' \
--data-urlencode 'sort=evtTime asc,id asc' \
--data-urlencode 'rows=1000' \
--data-urlencode 'cursorMark=*' \
--data-urlencode 'wt=json' | jq -r '.nextCursorMark'
응답의 nextCursorMark 를 다음 요청의 cursorMark 로 넘기고, 값이 직전과 같아지면 끝이다.
자주 쓰는 필드는 다음과 같다. 배포판과 Ranger 버전에 따라 구성이 다를 수 있으므로 실제 스키마를 확인한다 (확인 필요).
curl -s "${SOLR}/solr/ranger_audits/schema/fields?wt=json" | jq -r '.fields[].name'
| 필드 | 의미 |
|---|---|
evtTime |
이벤트 시각 |
reqUser |
요청한 사용자 |
repo · repoType |
대상 서비스(리포지터리) 이름과 종류 |
resource |
접근 대상 경로·테이블 |
access |
요청한 권한 |
result |
허용 1 · 거부 0 |
policy |
적용된 정책 ID |
cliIP |
클라이언트 IP |
거부된 접근만 뽑는 질의는 다음과 같다.
q=*:*&fq=result:0&fq=evtTime:[NOW-1DAY TO NOW]&sort=evtTime asc
주기 실행과 재시도, 오류 처리를 직접 만들지 않아도 되므로 NiFi 를 쓰는 편이 간단하다.
GenerateFlowFile(CRON) → InvokeHTTP(Solr select) → SplitJson($.response.docs[*])
→ PublishKafkaRecord
증분 수집을 하려면 마지막으로 읽은 시각을 상태로 남긴다. UpdateAttribute 의 상태 저장 기능이나 DistributedMapCache 에 last.evtTime 을 두고, 다음 질의의 fq 에 넣는다.
fq=evtTime:{${last.evtTime} TO NOW]
중괄호는 배타, 대괄호는 포함이다. 경계 값을 배타로 두어야 같은 건을 두 번 읽지 않는다. 다만 같은 시각에 여러 건이 있으면 경계에서 누락이 생길 수 있으므로, 중복을 허용하고 소비 측에서 id 로 걸러 내는 편이 안전하다.
전용 Solr 프로세서(GetSolr)도 있다. 마지막 수집 시각을 자체 상태로 관리하므로 증분 수집을 직접 만들지 않아도 된다. 다만 인증이 걸린 Solr 에 붙일 때 설정이 까다로울 수 있다.
Kafka 로 보낼 때는 파티셔닝 키를 정해 둔다. repo 나 reqUser 를 키로 두면 같은 대상의 기록이 같은 파티션으로 모여 순서가 유지된다.
파이썬으로 짧게 만들 수도 있다. 주기 실행은 스케줄러에 맡긴다.
import json
import requests
from kafka import KafkaProducer
SOLR = 'http://solr01.example.com:8983/solr/ranger_audits/select'
producer = KafkaProducer(
bootstrap_servers=['broker01.example.com:9092'],
value_serializer=lambda v: json.dumps(v).encode('utf-8'),
)
cursor = '*'
while True:
params = {
'q': '*:*',
'fq': 'evtTime:[NOW-1HOUR TO NOW]',
'sort': 'evtTime asc,id asc',
'rows': 1000,
'cursorMark': cursor,
'wt': 'json',
}
res = requests.get(SOLR, params=params, timeout=30).json()
for doc in res['response']['docs']:
producer.send('ranger-audit', doc)
nxt = res['nextCursorMark']
if nxt == cursor:
break
cursor = nxt
producer.flush()
Kerberos 가 걸린 Solr 이면 requests-kerberos 로 SPNEGO 인증을 붙인다.