지점·공장·엣지 데이터 수집

IT 엔지니어개발자

시작하기

기본 개념

여러 현장에서 생성되는 파일을 하나의 중앙 환경으로 모으기

지점의 업무 시스템, 공장의 생산 설비, 현장의 엣지 장비에서는 업무 자료와 생산 데이터, 로그와 결과 파일 등 다양한 파일이 생성됩니다.

지점·공장·엣지 데이터 수집은 각 현장의 파일 생성 위치를 중앙 수집 환경과 연결하고, 설정한 조건에 따라 본사 서버나 클라우드 스토리지로 파일을 자동 전송하는 방식입니다.

text
지점 A ──────┐
             │
공장 B ──────┼────→ 중앙 수집 환경 ────→ 본사 서버
             │              │
엣지 장비 C ─┘              └──────────→ 클라우드

각 현장의 파일을 하나의 중앙 수집 흐름으로 연결하면, 분산된 환경에서 생성되는 데이터를 지정된 저장 위치로 모아 이후 분석과 업무 처리에 활용할 수 있습니다.

수집 흐름

현장에서 생성된 파일을 감지해 중앙 저장 위치까지 자동으로 이어가기

현장에서 파일이 생성되거나 변경되면 설정한 수집 조건을 확인하고, 수집 대상 파일을 중앙 서버 또는 클라우드로 전송합니다.

text
현장 파일 생성
      │
      ▼
파일 변경 감지
      │
      ▼
수집 조건 확인
      │
      ▼
중앙 환경으로 전송
      │
      ▼
수집 결과 확인

파일 생성과 변경 시점 또는 정해진 일정에 따라 수집 작업을 시작하도록 구성할 수 있으며, 수집된 파일은 중앙 환경에서 분석과 후속 업무에 활용할 수 있습니다.

운영 효과

여러 현장의 데이터 흐름을 중앙에서 함께 관리하기

각 현장은 사용하는 장비와 파일 생성 위치, 수집 시점이 다를 수 있습니다. 이를 중앙 수집 환경에 연결하면 여러 현장에서 생성되는 파일을 하나의 흐름으로 관리할 수 있습니다.

현장 환경생성 파일중앙 수집 위치활용
지점업무 자료 · 보고서본사 서버업무 확인
공장생산 데이터 · 검사 결과클라우드분석 · 품질 관리
엣지 장비센서 데이터 · 로그중앙 분석 환경데이터 처리

IT 엔지니어

수집 환경

지점·공장·엣지 장비와 파일 생성 위치 연결하기

먼저 파일을 수집할 지점, 생산 설비와 엣지 장비를 중앙 관리 환경에 연결합니다.

각 장비에서 파일이 생성되는 폴더 또는 저장 위치를 지정하면 수집 작업에서 확인할 대상과 파일 경로를 구성할 수 있습니다.

text
Branch-A
└── /data/report

Factory-01
└── /production/result

Edge-Server-01
└── /logs/device

수집 경로

현장별 파일 위치를 중앙 저장 환경까지 연결하기

수집 대상 장비를 연결한 후에는 현장별 파일 경로와 중앙 저장 위치를 하나의 수집 경로로 구성합니다.

각 현장의 파일을 본사 서버로 모으거나, 데이터 유형과 활용 목적에 따라 클라우드 스토리지 또는 분석 환경으로 전송할 수 있습니다.

text
지점 A ────────┐
               │
공장 B ────────┼──→ 중앙 수집 ───→ 본사 서버
               │         │
엣지 장비 C ───┘         └───────→ Cloud Storage
수집 원본파일 경로중앙 저장 위치
지점 A/report/daily/data/branch
공장 B/production/result/data/factory
엣지 C/logs/device/data/edge

수집 조건

파일 이벤트와 일정에 따라 수집 작업 시작하기

파일 생성과 변경을 기준으로 수집 작업을 시작하거나, 정해진 일정에 따라 필요한 파일을 수집하도록 구성할 수 있습니다.

파일 경로와 유형을 기준으로 수집 대상을 지정하면 현장에서 생성되는 파일 중 필요한 데이터에 맞춰 수집 범위를 설정할 수 있습니다.

text
              수집 시작 조건
                    │
       ┌────────────┼────────────┐
       ▼            ▼            ▼
    파일 생성      파일 변경      정기 일정
       │            │            │
       └────────────┼────────────┘
                    ▼
                 파일 수집

수집 자동화

중앙 수집 이후 데이터 처리와 결과 저장까지 연결하기

파일 수집 작업은 여러 현장의 데이터를 중앙으로 모은 후 분석이나 변환, 별도 저장 작업까지 연결할 수 있습니다.

text
지점 데이터 ────┐
                │
공장 데이터 ────┼──→ 중앙 수집 ───→ 데이터 처리
                │                       │
엣지 데이터 ────┘                       ▼
                                    결과 저장

수집된 파일을 처리 장비와 연결하면 현장 데이터 수집 → 중앙 전송 → 데이터 처리 → 결과 저장까지 이어지는 자동화 흐름을 구성할 수 있습니다.

수집 현황

여러 현장의 수집 작업과 처리 결과를 중앙에서 확인하기

Runs와 Dataset 화면을 통해 여러 지점과 공장, 엣지 장비의 수집 작업을 함께 확인할 수 있습니다.

장비별 실행 상태와 진행률, 처리된 파일 수, 최근 실행 결과를 기준으로 현재 수집 현황을 관리합니다.

text
Branch-A
├── Status      Completed
├── Files       128
└── Last Run    Completed

Factory-01
├── Status      Running
├── Progress    68%
└── Last Run    In Progress

Edge-Server-01
├── Status      Completed
└── Files       2,431

이를 통해 여러 현장의 파일 수집 작업을 개별적으로 확인하는 과정 대신, 중앙 화면에서 전체 현황과 장비별 처리 상태를 함께 관리할 수 있습니다.

운영 대응

현장 연결과 파일 처리 상태를 확인하고 필요한 작업 다시 실행하기

수집 작업에서 확인이 필요한 경우 실행 상세 정보와 Activity Log를 통해 해당 현장의 장비 연결 상태와 파일 경로, 접근 범위, 처리 결과를 확인합니다.

text
수집 작업 실행
      │
      ▼
실행 상태 확인
      │
 ┌────┴─────┐
 ▼          ▼
Completed  확인 필요
 │          │
 ▼          ▼
결과 확인  장비 · 경로 · 파일 상태 확인
                │
                ▼
             설정 조정
                │
                ▼
             작업 재실행
                │
                ▼
              결과 확인

개발자

현장별 수집 자동화를 등록하고 연결 상태와 수집 현황 집계하기

연동 준비

공통 호출 코드와 경로 표기 준비하기

import os
import requests

BASE_URL = os.getenv("INNORIX_BASE_URL", "https://app.innorix.com").rstrip("/")
TOKEN = os.environ["INNORIX_ACCESS_TOKEN"]
WORKSPACE_ID = os.getenv("INNORIX_WORKSPACE_ID")   # optional; falls back to the current workspace

STATUS_COMPLETE = 2
TERMINAL = {2, 4, 5, 9, 99}          # complete / error / cancelled / partial / failed
NOT_SUCCEEDED = {4, 5, 9, 99}


def api(method, path, body=None, params=None):
    headers = {
        "Content-Type": "application/json",
        "Authorization": f"Bearer {TOKEN}",
    }

    if WORKSPACE_ID:
        headers["x-workspace-id"] = WORKSPACE_ID

    response = requests.request(
        method, BASE_URL + path,
        headers=headers, json=body, params=params, timeout=30,
    )

    payload = response.json() if response.content else {}

    if not response.ok:
        raise RuntimeError(payload.get("message") or f"HTTP {response.status_code}")

    return payload.get("data")


def is_terminal(detail):
    return detail.get("isTerminal", detail.get("status") in TERMINAL)
import base64
import time


def encode_path(device_id, raw_path):
    normalized = str(raw_path or "").replace("\\", "/")
    token = base64.b64encode(normalized.encode("utf-8")).decode("ascii")
    return f"{device_id}_ino_{token}"


def now_iso():
    return time.strftime("%Y-%m-%dT%H:%M:%S.000Z", time.gmtime())

전송 상태는 아래 값으로 판단합니다. 종료 상태는 다섯 개이고 성공에 해당하는 값은 완료(2)입니다.

상태 값의미종료
2완료
4오류
5취소
9부분 완료
99실패
1 · 6 · 12 · 13시작·전송 중·동기화 중·수신 중아니오

현장 목록 구성

현장 정보를 데이터로 두고 수집 작업 일괄 생성하기

현장이 수십 곳이면 화면에서 하나씩 만들기 어렵습니다. 현장 목록을 데이터로 두고 순회합니다.

SITES = [
    {"device": "branch-seoul", "path": "/report/daily", "target": "/data/branch"},
    {"device": "branch-busan", "path": "/report/daily", "target": "/data/branch"},
    {"device": "factory-01", "path": "/production/result", "target": "/data/factory"},
    {"device": "edge-line-01", "path": "/logs/device", "target": "/data/edge"},
]

CENTRAL = "device-hq-01"


def site_target(site):
    # sites reuse the same file names, so put the site id in the target path
    return f"{site['target']}/{site['device']}"

도착 경로를 나누지 않으면 현장마다 result.csv처럼 이름이 같은 파일이 서로 덮어씁니다.

수집 자동화 등록

현장별로 수집 작업을 만들고 실행 조건 정하기

def build_collection(site, central, schedule=None, options=None):
    name = f"collect {site['device']}"
    sync = schedule is None

    body = {
        "name": name,
        "flowName": name,
        "transferType": "sync" if sync else "normal",
        "timezone": "Asia/Seoul",
        "step": 1,
        "isUpcoming": False,
        "details": [
            {
                "senderId": site["device"],
                "receiverId": central,
                "sourceItem": [
                    {
                        "hash": encode_path(site["device"], site["path"]),
                        "filePath": site["path"],
                        "isDir": True,
                    }
                ],
                "targetPath": encode_path(central, site_target(site)),
                "step": 1,
                "transferOptions": {
                    "noSchedule": sync,
                    "target-action": "numbering",
                    "send-fileoption": {},
                    **({"syncType": 1} if sync else {}),
                    **(options or {}),
                },
            }
        ],
        "schedules": [schedule or {
            "type": "none", "startDateType": "now",
            "startDate": now_iso(), "timezone": "Asia/Seoul",
        }],
    }

    return body


DAILY = {
    "type": "day", "startDateType": "now",
    "hour": "02", "minute": "00", "ampm": "am",
    "startDate": now_iso(), "timezone": "Asia/Seoul",
}

collectors = {
    site["device"]: api("POST", "/api/automations",
                        build_collection(site, CENTRAL, DAILY))["automationId"]
    for site in SITES
}

자동화 요청에서 반드시 지켜야 하는 항목이 네 개 있습니다.

항목지정 방식
isUpcoming반드시 false. 서버 기본값 true는 요청에 담긴 일정을 무시하고 5분짜리 일회성 일정으로 대체합니다. triggerAutomation이 붙은 단계는 서버가 false로 강제하므로, 트리거가 없는 첫 단계에만 직접 지정하면 됩니다
step최상위와 details 양쪽에 넣습니다. 흐름 안에서의 홉 위치입니다
sourceItemhash(경로 토큰)와 filePath(평문 경로)를 함께 넣습니다
syncTypetransferOptions 안에 넣습니다. 1은 단방향, 2는 양방향입니다

네 항목 모두 누락해도 등록은 성공하고 실행 시점에 동작이 달라집니다. 반복 일정을 등록했는데 한 번만 실행되고 끝났다면 isUpcoming부터 확인합니다.

실행 방식구성적합한 현장
일정 실행transferType: normal + 일정파일이 정해진 시각에 만들어지는 지점
생성 감지transferType: sync + transferOptions.syncType파일이 수시로 생기는 설비와 엣지 장비

현장마다 파일 생성 시점이 다르므로 한 방식으로 통일하지 않아도 됩니다.

중복 등록 방지

같은 수집 작업이 두 번 만들어지지 않게 하기

같은 이름의 자동화가 있어도 새로 만들어집니다. 현장 목록을 다시 순회하면 수집이 두 번 실행됩니다.

def find_automation(name):
    # the name we send is stored as flowName in the response
    # automationName is a server generated id like T4037-8500-1815, not the name we set.
    for page in range(1, 6):
        result = api("GET", "/api/automations",
                     params={"page": page, "size": 100, "search": name}) or {}

        items = [item
                 for flow in result.get("automations") or []
                 for item in flow.get("automations") or []]

        for item in items:
            if item.get("flowName") == name:
                return item

        if len(items) < 100:
            return None

    return None


def ensure_collection(site, central, schedule=None):
    name = f"collect {site['device']}"

    if find_automation(name):
        return None

    return api("POST", "/api/automations",
               build_collection(site, central, schedule))["automationId"]

자동화 목록은 흐름 그룹으로 중첩되어 반환되므로 안쪽 배열까지 순회해야 합니다. 서버 검색은 부분 일치이므로, 받은 결과에서 이름이 정확히 일치하는 항목만 골라냅니다.

오프라인 현장 대응

연결이 끊긴 현장을 찾고 복구 후 밀린 파일 보내기

현장 장비는 네트워크가 불안정한 경우가 있습니다. 수집 실패인지 장비 연결 해제인지에 따라 대응이 달라집니다.

def site_state(device_id):
    state = api("GET", f"/api/devices/{device_id}/connectivity") or {}
    return bool(state.get("isConnected")), state.get("stateLabel")


for site in SITES:
    connected, label = site_state(site["device"])

    if not connected:
        print(f"{site['device']:20} {label}")

응답의 연결 여부는 isConnected입니다. 함께 오는 stateLabel은 화면에 그대로 쓸 수 있습니다.

연결이 돌아오면 그동안 쌓인 파일을 한 번에 보냅니다. 소스와 대상 목록을 비교해 누락된 파일만 선별합니다.

def list_files(device_id, path):
    found, page = {}, 1

    while True:
        result = api("GET", f"/api/devices/{device_id}/files", params={
            "path": path, "page": page, "size": 200, "type": "file",
        }) or {}

        for item in result.get("items") or []:
            found[item["name"]] = item.get("size")

        if page >= (result.get("lastPage") or 1):
            return found

        page += 1


def catch_up(site, central):
    source_files = list_files(site["device"], site["path"])
    missing = sorted(set(source_files) - set(list_files(central, site_target(site))))

    if not missing:
        return None

    # send file lists through sourceItem; sourcePaths treats every path as a folder
    return api("POST", "/api/transfers/manual", {
        "sourceDevice": site["device"],
        "targetDevice": central,
        "targetPath": site_target(site),
        "sourceItem": [
            {"path": f"{site['path'].rstrip('/')}/{name}",
             "isDir": False,
             "fileSize": source_files[name]}
            for name in missing
        ],
        "sendAllFolder": False,
        "transferOptions": {"target-action": "numbering"},
    })["monitorId"]

파일 크기를 함께 넘기면 서버가 항목별 크기를 다시 조회하지 않아 빨라집니다.

도착 정책이 numbering이면 이름에 번호가 붙어 쌓입니다. 그 구성에서는 이름 비교 대신 전송 이력에 없는 파일을 기준으로 찾습니다.

수집 현황 집계

여러 현장의 수집 결과를 한 번에 확인하기

def paginate(path, params=None, limit=200, max_pages=50):
    query = dict(params or {})
    query["limit"] = limit
    cursor = None

    for _ in range(max_pages):
        if cursor:
            query["cursor"] = cursor

        result = api("GET", path, params=query) or {}

        for record in result.get("data") or []:
            yield record

        pagination = result.get("pagination") or {}

        if not pagination.get("hasMore"):
            return

        cursor = pagination.get("nextCursor")

        if not cursor:
            return
from collections import Counter
from datetime import datetime, timedelta, timezone


def site_summary(device_id, days=1):
    end = datetime.now(timezone.utc)
    fmt = "%Y-%m-%dT%H:%M:%SZ"

    rows = list(paginate(f"/api/devices/{device_id}/transfer-history", params={
        "startDate": (end - timedelta(days=days)).strftime(fmt),
        "endDate": end.strftime(fmt),
    }))

    failed = [r for r in rows if r.get("status") in NOT_SUCCEEDED]
    return {"total": len(rows), "failed": len(failed)}


header = "{:20} {:^6} {:>6} {:>6}".format("site", "up", "collect", "fail")
print(header)
print("-" * len(header))

for site in SITES:
    connected, _ = site_state(site["device"])
    summary = site_summary(site["device"])

    print(f"{site['device']:20} {'O' if connected else 'X':^6}"
          f" {summary['total']:>6} {summary['failed']:>6}")

전송 이력은 data.data 배열로, 페이징 정보는 data.pagination으로 반환됩니다.

예외 대응

실패한 수집을 확인하고 다시 실행하기

def wait(monitor_id, timeout=3600, interval=3):
    deadline = time.time() + timeout

    while time.time() < deadline:
        detail = api("GET", f"/api/transfers/{monitor_id}")

        if is_terminal(detail):
            return detail

        time.sleep(interval)

    raise TimeoutError(monitor_id)


def failed_files(monitor_id):
    result = api("GET", f"/api/transfers/{monitor_id}/files", params={
        "state": "any", "size": 500,
    }) or {}

    return [r for r in (result.get("children") or [])
            if r.get("status") in NOT_SUCCEEDED]


def retry_failed(monitor_id):
    rows = failed_files(monitor_id)

    if not rows:
        return 0

    api("POST", f"/api/transfers/{monitor_id}/retry", {
        "filesRetry": [
            {"filePath": r["sourceFilePath"], "isDir": bool(r.get("isFolder"))}
            for r in rows
        ]
    })

    return len(rows)
for device, automation_id in collectors.items():
    runs = api("GET", f"/api/automations/{automation_id}/executions") or []

    if not runs:
        print(f"{device}: no run history - check registration and start condition")
        continue

    latest = runs[0]

    if latest["status"] != STATUS_COMPLETE:
        print(f"{device}: retried {retry_failed(latest['monitorId'])} files")

실행 이력은 최신 회차가 배열 앞에 옵니다. 이력이 비어 있으면 자동화가 등록만 되고 실행되지 않은 것이므로 isUpcoming과 시작 조건을 확인합니다.

확인 항목확인 내용
현장 장비isConnectedstateLabel
수집 작업현장별 자동화와 실행 조건
도착 경로현장 식별자로 구분된 저장 위치
수집 현황현장별 수집 건수와 실패
밀린 파일오프라인 기간의 누락 파일