수집기는 장비가 말하는 것(레지스터, 파일, HTTP 응답)을 레코드(구조화된 값)나 블롭(파일)으로 바꿔 로컬 큐에 넣습니다. 중앙으로 보내는 일은 큐와 업로더가 하므로, 수집기는 링크가 끊겨 있어도 똑같이 돕니다.

edge.yamlcollectors: 에 선언합니다. 코드를 쓰지 않습니다.

종류엑스트라읽는 것내는 것화면에서
http없음다른 서버의 JSON 엔드포인트폴링마다 레코드 1건(또는 항목마다)수집 데이터
push없음없음(로컬 API 로 들어옴)PUT 마다 레코드 1건수집 데이터
modbusmodbusPLC 레지스터 맵(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_ms1000폴링 주기
timeout_s10요청 제한 시간
verify_tlstrueTLS 검증
headers{}요청 헤더
auth없음HTTP Basic 인증용 {username, password}
kindhttp레코드 종류 이름
ts_field""응답(또는 항목)에서 관측 시각을 읽을 필드
items_path""응답 안 배열 위치(data.items, 본문 자체가 배열이면 .)
id_field""항목의 고유 키. items_path 를 쓰면 필수
max_body_bytes1MiB응답 크기 상한

레코드 내용은 {"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_pathid_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 는 선택입니다. 없으면 받은 시각을 씁니다.
  • kindpriorityYAML 이 정합니다. 보내는 쪽이 아무도 모르는 종류 이름으로 데이터를 넣을 수 없게 하려는 것입니다.
  • 선언된 입구는 수집기 목록에 마지막 수신 시각과 함께 보입니다. 보내는 쪽이 멈추면 오류조차 오지 않습니다. 이 시각이 더 이상 바뀌지 않는 것이 유일한 신호입니다.

선언 없이 넣기

설정에 없는 프로세스를 위해 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 안에 바로 적어도 됨)
옵션기본값설명
host127.0.0.1PLC 주소
port502포트
unit_id1기본 유닛 ID
interval_ms1000폴링 주기
schema (또는 register_map)(필수)레지스터 맵(파일 경로 또는 YAML 안에 바로 쓴 매핑)
kindmodbus레코드 종류 이름

레지스터 맵은 평범한 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*파일 이름 패턴
kindblob블롭 종류 이름
interval_s2.0훑는 주기
recursivefalse하위 폴더까지
delete_aftertrue스풀로 옮긴 뒤 원본 삭제. false 면 복사만 하고 같은 파일은 다시 집지 않음
stable_checks1크기가 몇 번 연속 같아야 다 쓴 파일로 볼지

아직 쓰는 중인 파일을 집지 않도록, 크기가 한 주기 동안 그대로인 파일만 가져갑니다. 에이전트 계정이 그 폴더에 쓰기 권한도 가져야 합니다. 원본을 지우지 못하면 같은 파일을 반복해 집지는 않지만 수집기에 오류로 표시됩니다.

커스텀 수집기

기본 네 가지로 안 되는 장비(시리얼 포트, 전용 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} 같은 레코드로 나타납니다.

2026-09-21 기준 플랫폼에 맞춰 작성했습니다.

© Geo-MLOps