데이터 수집기
http · push · modbus · watchdir 수집기 설정과 커스텀 수집기 만들기
수집기는 장비가 말하는 것(레지스터, 파일, HTTP 응답)을 레코드(구조화된 값)나 블롭(파일)으로 바꿔 로컬 큐에 넣습니다. 중앙으로 보내는 일은 큐와 업로더가 하므로, 수집기는 링크가 끊겨 있어도 똑같이 돕니다.
edge.yaml 의 collectors: 에 선언합니다. 코드를 쓰지 않습니다.
| 종류 | 엑스트라 | 읽는 것 | 내는 것 | 화면에서 |
|---|---|---|---|---|
http | 없음 | 다른 서버의 JSON 엔드포인트 | 폴링마다 레코드 1건(또는 항목마다) | 수집 데이터 |
push | 없음 | 없음(로컬 API 로 들어옴) | PUT 마다 레코드 1건 | 수집 데이터 |
modbus | modbus | PLC 레지스터 맵(Modbus TCP) | 폴링마다 레코드 1건 | 수집 데이터 |
watchdir | 없음 | 다른 프로그램이 폴더에 쓴 파일 | 파일마다 블롭 1개 | 파일 |
공통 필드
collectors:
- type: http # 종류 (필수)
name: gateway-1 # 이름. 화면의 수집기 목록과 레코드의 collector 값
enabled: true # false 면 건너뜀
priority: 50 # 0~100. 기본 50
options: { ... } # 종류별 옵션
priority 는 보존과 업로드가 모두 읽습니다. 한도를 넘으면 낮은 것부터 버리고, 보낼 때는 높은 것부터 보냅니다. sync.urgent_priority(기본 90) 이상은 전송 시간대 제한도 무시합니다. 특별히 중요한 것이 아니면 50 그대로 둡니다.
설정이 잘못된 수집기 하나는 로그에 남기고 건너뜁니다. 나머지 수집기와 중앙 통신은 계속 돕니다. 수집기 상태(running·disconnected·stopped)와 마지막 오류는 디바이스 상세의 수집기 카드와 geo-mlops-edge status 에 나옵니다.
옵션은 options: 아래에 써도 되고, 같은 줄에 바로 써도 됩니다(둘 다 받습니다).
http: 다른 서버를 폴링
- type: http
name: gateway-1
options:
url: http://10.0.0.7/api/current
interval_ms: 1000
timeout_s: 10
verify_tls: true
headers: { Authorization: "Bearer <your-secret>" } # 선택
auth: { username: edge, password: <your-secret> } # 선택, HTTP Basic
ts_field: measured_at # 응답에 관측 시각이 있으면 그 필드
max_body_bytes: 1MiB # 응답 크기 상한. 0 이면 검사 안 함
| 옵션 | 기본값 | 설명 |
|---|---|---|
url | (필수) | JSON 을 돌려주는 주소 |
interval_ms | 1000 | 폴링 주기 |
timeout_s | 10 | 요청 제한 시간 |
verify_tls | true | TLS 검증 |
headers | {} | 요청 헤더 |
auth | 없음 | HTTP Basic 인증용 {username, password} |
kind | http | 레코드 종류 이름 |
ts_field | "" | 응답(또는 항목)에서 관측 시각을 읽을 필드 |
items_path | "" | 응답 안 배열 위치(data.items, 본문 자체가 배열이면 .) |
id_field | "" | 항목의 고유 키. items_path 를 쓰면 필수 |
max_body_bytes | 1MiB | 응답 크기 상한 |
레코드 내용은 {"collector": "gateway-1", "body": <응답 JSON>} 입니다.
같은 값이 와도 매번 한 건입니다. "12:00:01 에도 21.5 였다"는 것도 데이터이기 때문입니다. 그래서 데이터 양은 주기 × 응답 크기 로 가늠할 수 있습니다. 2 KiB 를 1초마다 읽으면 하루 약 177 MB 입니다. 기본 보존 한도(50 GiB / 30일)에서는 기간 한도가 먼저 걸립니다.
지난 기록 목록을 돌려주는 엔드포인트
"최근 알람 N건"처럼 지난 기록 목록을 돌려주는 엔드포인트는 폴링할 때마다 이미 받은 행이 다시 옵니다. 배열 위치와 키를 알려 주면 항목마다 한 번씩만 큐에 넣습니다.
- type: http
name: alarms
priority: 70
options:
url: http://10.0.0.7/api/alarms
interval_ms: 5000
items_path: data.items
id_field: id
판단 기준은 하나입니다. 이 응답은 지금 상태인가, 지난 기록 목록인가? 지금 상태라면 폴링마다 한 건씩 쌓으면 됩니다. 지난 기록 목록이라면 items_path 와 id_field 로 키를 알려 줘야 합니다.
상대 서버가 멈췄거나, 느리거나, 이상한 값을 주는 일은 흔히 있는 상황으로 봅니다. 수집기는 오류를 기록하고 disconnected 로 표시합니다. 그런 다음 1초부터 30초까지 간격을 늘려 가며 다시 시도해 스스로 회복합니다.
push: 밖에서 밀어 넣는 입구
로봇 컨트롤러·비전 PC·다른 언어로 된 서비스가 LAN 에서 에이전트로 레코드를 밀어 넣을 때 씁니다.
- type: push
name: robot-1
priority: 60
options:
kind: robot
보내는 쪽은 자기가 정한 id 로 PUT 합니다.
curl -X PUT http://edge-pc:8600/api/v1/collectors/robot-1/records/evt-1 \
-H 'content-type: application/json' \
-d '{"payload": {"step": 3}, "ts": "2026-09-16T01:02:03Z"}'
201은 저장했다는 뜻이고,200과"duplicate": true는 이미 받은 id 라는 뜻입니다. 응답을 못 받아 다시 보내도 한 건만 남습니다.ts는 선택입니다. 없으면 받은 시각을 씁니다.kind와priority는 YAML 이 정합니다. 보내는 쪽이 아무도 모르는 종류 이름으로 데이터를 넣을 수 없게 하려는 것입니다.- 선언된 입구는 수집기 목록에 마지막 수신 시각과 함께 보입니다. 보내는 쪽이 멈추면 오류조차 오지 않습니다. 이 시각이 더 이상 바뀌지 않는 것이 유일한 신호입니다.
선언 없이 넣기
설정에 없는 프로세스를 위해 POST /api/v1/records, POST /api/v1/blobs 도 열려 있습니다. 대신 id 를 에이전트가 만들므로 재시도하면 두 건이 되고, 수집기 목록에도 나오지 않습니다.
| 입구 | id | 재시도 안전 | kind 를 정하는 쪽 | 플릿에 보임 |
|---|---|---|---|---|
| 수집기 | 에이전트가 생성 | 해당 없음 | YAML | 예 |
POST /records, POST /blobs | 에이전트가 생성 | 아니오 | 보내는 쪽 | 아니오 |
PUT /collectors/{name}/records/{id} | 보내는 쪽 | 예 | YAML | 예 |
modbus: PLC 레지스터 읽기
pip install 'geo-mlops-sdk[edge,modbus]' 가 필요합니다.
- type: modbus
name: line-1
priority: 50
options:
host: 10.0.0.5
port: 502
unit_id: 1
interval_ms: 1000
schema: /etc/geo-mlops/plc.yaml # 레지스터 맵 파일 (YAML 안에 바로 적어도 됨)
| 옵션 | 기본값 | 설명 |
|---|---|---|
host | 127.0.0.1 | PLC 주소 |
port | 502 | 포트 |
unit_id | 1 | 기본 유닛 ID |
interval_ms | 1000 | 폴링 주기 |
schema (또는 register_map) | (필수) | 레지스터 맵(파일 경로 또는 YAML 안에 바로 쓴 매핑) |
kind | modbus | 레코드 종류 이름 |
레지스터 맵은 평범한 YAML 입니다.
# /etc/geo-mlops/plc.yaml
version: "1"
unit_id: 1
fields:
# type: holding | input | coil | discrete
# dtype: bool | uint16 | int16 | uint32 | int32 | float32
- { name: temperature.zone1, type: holding, address: 100, dtype: uint16, scale: 0.1 }
- { name: temperature.zone2, type: holding, address: 101, dtype: uint16, scale: 0.1 }
# 32비트 값은 레지스터 두 개. word_order 가 틀리면 오류 대신 그럴듯한 엉터리 값이 나온다
- { name: flow_rate, type: holding, address: 110, dtype: float32, word_order: big }
- { name: cycle_count, type: holding, address: 112, dtype: uint32 }
- { name: running, type: coil, address: 5 }
- { name: fault, type: discrete, address: 12 }
연속된 주소는 한 번에 읽습니다. 레코드 내용은 {"collector": "line-1", "version": "1", "fields": {"temperature.zone1": 21.3, ...}} 입니다.
PLC 가 꺼져 있으면 disconnected 로 표시하고, 1초부터 30초까지 간격을 늘려 가며 다시 접속합니다.
watchdir: 폴더에 새로 생긴 파일
카메라, 라이다, 오래된 장비 프로그램이 폴더에 쓰는 파일을 가져와 올립니다. 파일은 파일 탭으로 갑니다.
- type: watchdir
name: cam-0
priority: 20
options:
path: /data/incoming
pattern: "*.jpg"
interval_s: 2
delete_after: true
| 옵션 | 기본값 | 설명 |
|---|---|---|
path | . | 감시할 폴더(없으면 만듦) |
pattern | * | 파일 이름 패턴 |
kind | blob | 블롭 종류 이름 |
interval_s | 2.0 | 훑는 주기 |
recursive | false | 하위 폴더까지 |
delete_after | true | 스풀로 옮긴 뒤 원본 삭제. false 면 복사만 하고 같은 파일은 다시 집지 않음 |
stable_checks | 1 | 크기가 몇 번 연속 같아야 다 쓴 파일로 볼지 |
아직 쓰는 중인 파일을 집지 않도록, 크기가 한 주기 동안 그대로인 파일만 가져갑니다. 에이전트 계정이 그 폴더에 쓰기 권한도 가져야 합니다. 원본을 지우지 못하면 같은 파일을 반복해 집지는 않지만 수집기에 오류로 표시됩니다.
커스텀 수집기
기본 네 가지로 안 되는 장비(시리얼 포트, 전용 SDK 등)는 수집기를 직접 만들어 등록합니다. register_collector 는 같은 프로세스 안에서 에이전트를 띄우기 전에 불러야 하므로, geo-mlops-edge run 대신 작은 실행 스크립트를 씁니다.
# my_edge.py
import asyncio
import sys
from datetime import datetime, timezone
from pathlib import Path
from geo_mlops_sdk.edge.collectors import CollectorBase, register_collector
from geo_mlops_sdk.edge.daemon import run
from geo_mlops_sdk.edge.settings import EdgeSettings
class CounterCollector(CollectorBase):
type_name = "counter" # 화면에 보이는 종류
def __init__(self, name, *, priority=50, interval_s=5.0):
super().__init__(name, priority=priority)
self.interval_s = interval_s
self._task = None
@classmethod
def from_options(cls, *, name, priority, options):
return cls(name, priority=priority,
interval_s=float(options.get("interval_s", 5)))
async def start(self, sink):
self.state = "running"
self._task = asyncio.create_task(self._loop(sink))
async def stop(self):
if self._task:
self._task.cancel()
self.state = "stopped"
async def _loop(self, sink):
value = 0
while True:
value += 1
now = datetime.now(timezone.utc)
await sink.record("counter", {"value": value},
priority=self.priority, ts=now)
self.note_emit(now) # 수집기 카드의 '마지막 수신' 갱신
await asyncio.sleep(self.interval_s)
register_collector("counter", CounterCollector.from_options)
if __name__ == "__main__":
config = Path(sys.argv[1]) if len(sys.argv) > 1 else None
sys.exit(run(EdgeSettings.load(config)))
collectors:
- type: counter
name: counter-1
options:
interval_s: 15
python my_edge.py /etc/geo-mlops/edge.yaml
- 팩토리는
factory(name=..., priority=..., options=...)로 불립니다. sink.record(kind, payload, priority=, ts=, meta=, record_id=)는 레코드를,sink.blob(kind, 경로|bytes, filename=, priority=, move=)는 파일을 큐에 넣습니다.record_id를 주면 같은 id 는 한 번만 들어갑니다.CollectorBase를 상속하지 않아도name,start(sink),stop(),status()만 있으면 됩니다. 상속하면note_emit()·note_error()로 화면의 수집기 상태가 맞게 채워집니다.run()이 돌려주는 종료 코드(재시작 요청이면 3)를 그대로sys.exit에 넘겨야 systemd 가 올바르게 다시 띄웁니다.
위 예시는 실제로 돌려 확인했습니다. 디바이스 상세 수집기 카드에 counter-1 (counter) running 으로, 수집 데이터 탭에 {"value": 8} 같은 레코드로 나타납니다.