미디어·데이터 자동화

IT 엔지니어개발자

시작하기

기본 개념

파일을 수집해 처리와 결과 전달까지 자동으로 연결하기

미디어와 데이터 작업에서는 파일이 생성되는 위치와 처리하는 장비, 결과를 활용하는 위치가 서로 다를 수 있습니다.

예를 들어 여러 촬영 장비에서 생성된 영상 파일을 수집해 변환 서버로 전송하고, 처리된 결과를 스토리지에 저장한 후 다음 업무 환경으로 전송하는 흐름을 구성할 수 있습니다.

text
파일 생성
    │
    ▼
파일 수집
    │
    ▼
파일 처리
    │
    ▼
결과 활용

각 단계의 결과는 다음 작업의 입력으로 이어지며, 파일이 생성된 이후 수집과 처리, 결과 활용까지 하나의 흐름으로 구성할 수 있습니다.

자동화 흐름

여러 장비의 파일을 수집해 정해진 순서로 처리하기

자동화 흐름에서는 여러 장비에서 생성된 파일을 하나의 작업으로 수집하고, 수집 결과를 필요한 처리 장비로 연결합니다.

처리가 완료되면 생성된 결과 파일을 지정된 위치로 전송해 다음 업무에서 활용할 수 있습니다.

text
Source A ───┐
            │
Source B ───┼──→ Collect
            │       │
Source C ───┘       ▼
                  Process
                     │
                     ▼
                  Result

파일 상태와 작업 완료 조건을 기준으로 다음 단계가 실행되도록 설정하면 각 파일 작업을 순서에 따라 자동으로 연결할 수 있습니다.

자동화 효과

단계마다 반복되는 파일 이동과 관리 작업 줄이기

여러 단계의 파일 작업을 개별적으로 진행하면 파일 준비와 전송, 처리 결과 확인이 반복됩니다.

자동화 흐름을 구성하면 앞선 단계의 결과를 다음 작업과 연결해 파일이 생성된 이후의 처리 과정을 하나의 흐름으로 관리할 수 있습니다.

구분개별 파일 처리자동화 흐름
파일 수집장비별 파일을 각각 확인여러 장비의 파일을 하나의 흐름으로 수집
처리 실행파일 준비 후 작업 실행수집 결과에 따라 다음 작업 실행
결과 활용완료 파일을 다음 위치로 관리결과를 지정된 위치와 업무로 연결
진행 확인단계별 작업을 각각 확인전체 흐름과 단계별 결과 확인

이를 통해 여러 장비와 시스템에 분산된 파일 작업을 하나의 자동화 흐름으로 연결할 수 있습니다.

IT 엔지니어

미디어와 데이터 파일의 자동 처리 흐름 구성하기

수집 환경

파일이 생성되는 장비와 수집 위치 연결하기

먼저 파일이 생성되는 장비와 수집할 파일 위치를 연결합니다.

촬영 장비, 업무 서버, 데이터 수집 장비, 스토리지 등 파일이 생성되는 환경을 Source로 구성하고, 각 장비에서 사용할 파일 경로를 지정합니다.

text
Source Devices
      │
 ┌────┼────┐
 ▼    ▼    ▼
Cam  Server Storage
 │      │      │
 └──────┼──────┘
        ▼
    Collection

수집 환경을 구성하면 여러 위치에서 생성되는 파일을 하나의 자동화 흐름으로 연결할 수 있습니다.

수집한 파일을 필요한 처리 장비로 이어가기

수집된 파일은 변환 서버, 분석 서버, AI 처리 장비 등 필요한 작업을 수행하는 환경으로 연결합니다.

여러 Source의 파일을 하나의 처리 장비로 모으거나, 파일 종류와 처리 조건에 따라 각각 다른 작업으로 연결할 수 있습니다.

text
Source A ───┐
            ├──→ Collect ───→ Process Server
Source B ───┤                       │
            │                       ▼
Source C ───┘                   Processing

처리 작업이 완료되면 생성된 파일을 다음 결과 단계로 연결합니다.

결과 전달

처리 결과를 저장하고 다음 업무 환경으로 연결하기

처리된 결과 파일은 스토리지, 서버, 애플리케이션 등 다음 업무에서 사용할 위치로 전송합니다.

결과를 저장하는 작업과 다음 업무 환경으로 전송하는 작업을 하나의 결과 단계로 연결하면, 처리 완료 후 파일 활용까지 자동으로 이어갈 수 있습니다.

text
Processing
     │
     ▼
Result Files
     │
 ┌───┴───────┐
 ▼           ▼
Storage   Next System

필요한 경우 결과 전달 이후 알림이나 추가 처리 작업을 연결해 파일 활용 흐름을 확장할 수 있습니다.

실행 조건

파일과 작업 상태에 따라 다음 단계 실행하기

각 단계는 지정된 조건에 따라 실행되도록 구성할 수 있습니다.

예를 들어 파일이 생성되면 수집 작업을 시작하고, 수집이 완료되면 처리 작업을 실행하며, 처리 결과가 준비되면 결과 전달 작업을 시작하도록 연결할 수 있습니다.

text
File Ready
    │
    ▼
Collect Complete?
    │
   Yes
    │
    ▼
Start Processing
    │
    ▼
Result Ready?
    │
   Yes
    │
    ▼
Deliver Result

파일 상태와 앞선 작업의 결과를 실행 조건으로 활용하면 여러 작업을 정해진 순서로 연결할 수 있습니다.

결과 확인

전체 흐름과 단계별 처리 상태 함께 확인하기

자동화 흐름이 실행되면 전체 워크플로의 상태와 각 단계의 처리 결과를 함께 확인합니다.

수집된 파일과 처리 작업, 결과 전달 상태를 기준으로 현재 작업이 어느 단계까지 진행되었는지 확인할 수 있습니다.

text
Workflow Run
     │
 ┌───┼──────────────┐
 ▼   ▼       ▼      ▼
Collect Process Result Complete
  ✓      ✓       ●
                 │
              Running

확인 단계주요 확인 내용
수집수집된 파일과 실행 상태
처리처리 작업과 진행 결과
결과저장 또는 전송된 결과 파일
전체 실행워크플로 진행 상태와 실행 시간

전체 흐름과 개별 단계의 상태를 함께 확인하면 현재 작업의 진행 상황과 처리 결과를 빠르게 파악할 수 있습니다.

예외 대응

확인이 필요한 단계를 확인하고 작업 다시 실행하기

실행 과정에서 추가 확인이 필요한 경우 해당 단계의 상세 정보와 실행 기록을 확인합니다.

파일 수집, 처리 장비, 결과 전달 중 확인이 필요한 작업을 선택하고 장비 연결, 파일 경로, 처리 상태를 점검한 후 필요한 작업을 다시 실행할 수 있습니다.

text
Workflow Run
     │
     ▼
Status Check
     │
 ┌───┴────────┐
 ▼            ▼
Completed   Attention
                │
                ▼
          View Details
                │
                ▼
          Check Settings
                │
                ▼
              Retry
                │
                ▼
          Result Confirmed

확인 항목확인 내용후속 작업
수집 환경장비 연결과 파일 경로환경 확인 후 재실행
처리 작업처리 상태와 실행 결과처리 환경 확인
결과 전달대상 위치와 전송 상태연결 상태 확인 후 재실행
실행 기록단계별 작업 정보상세 내용 확인 후 조치

개발자

확장자와 크기로 대상을 고른 뒤 처리 장비로 모으고 결과를 다음 장비로 넘기기

연동 준비

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

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시작·전송 중·동기화 중·수신 중아니오

대상 선별

확장자·크기·이름 조건으로 보낼 파일 걸러내기

소스 폴더 전체가 아니라 조건에 맞는 파일만 보내도록 전송 옵션에 필터를 지정합니다.

def build_filter(exts=None, min_size=None, exclude=None):
    file_option = {}

    if exts:
        # extension whitelist, without the leading dot
        file_option["extension"] = {
            "extension": [e.lstrip(".").lower() for e in exts],
            "allow": True,
        }

    if min_size is not None:
        # over and equal both True means size or larger
        file_option["fileSize"] = {"size": min_size, "over": True, "equal": True}

    if exclude:
        # allow=False excludes files whose name contains this. Server matching is case sensitive.
        file_option["fileName"] = {"name": exclude, "allow": False}

    return {"send-fileoption": file_option} if file_option else {}

확장자 필터는 send-fileoption.extension을 씁니다. send-filetype-cus 정규식은 확장자를 뗀 파일명에만 매칭되므로 확장자 조건으로는 동작하지 않습니다.

필터위치동작
확장자send-fileoption.extensionallow: true 면 이 확장자만 전송
크기send-fileoption.fileSizeover·equal 로 이상·이하 지정
이름send-fileoption.fileNameallow: false 면 포함된 파일 제외

여러 필터를 함께 주면 AND로 결합됩니다. 모두 통과한 파일만 전송됩니다.

FILTER = build_filter(exts=["mp4", "mov", "wav"], min_size=1024, exclude=".tmp")

조건이 의도대로 걸리는지는 자동화에 넣기 전에 검색으로 확인합니다.

page = api("POST", f"/api/devices/{SOURCE_ID}/files/search",
           {"path": SOURCE_PATH, "pageSize": 500})

matched = [i for i in page["items"]
           if i["type"] == "file" and i["name"].lower().endswith((".mp4", ".mov"))]

print(f"{len(matched)} matched")

검색은 조건 파라미터를 받지 않으므로 받은 결과에서 판단합니다.

파일 수집

여러 장비의 파일을 처리 장비 한 곳으로 모으기

하나의 전송은 소스 장비 하나를 다룹니다. 장비마다 전송을 만들고 도착 경로를 소스별로 나눕니다.

SOURCES = [
    ("device-cam-01", "/media/raw"),
    ("device-cam-02", "/media/raw"),
    ("device-mic-01", "/audio/raw"),
]

for source, source_path in SOURCES:
    # give each source its own folder so file names do not collide
    target_path = f"/work/incoming/{source}"

    monitor_id = api("POST", "/api/transfers/manual", {
        "sourceDevice": source,
        "targetDevice": PROCESS_ID,
        "targetPath": target_path,
        "sourcePaths": [source_path],
        "sendAllFolder": True,
        "transferOptions": {"target-action": "numbering", **FILTER},
    })["monitorId"]

    print(source, "->", target_path, monitor_id)

장비별로 나뉘어 있으면 한 장비가 오프라인이어도 나머지 수집은 그대로 진행됩니다. 도착 경로를 나누지 않으면 장비마다 같은 이름의 파일이 서로 덮어씁니다.

저장 위치 제어

도착 경로 아래의 폴더 구조 정하기

TARGET_OPTIONS = {
    "savepath": True,          # keep the source folder structure (lowercase p)
    "optionPath": 3,           # how many trailing path segments to keep
    "target-action": "numbering",
}
필드타입내용
savepathboolean소스의 폴더 구조를 유지할지 여부
optionPathinteger유지할 소스 경로의 뒤쪽 단계 수
target-actionstring같은 이름이 있을 때의 처리 정책

여러 장비에서 같은 이름의 파일이 오는 구성이라면 optionPath 값을 늘려 출처를 구분합니다. 날짜별로 나눠 담으려면 옵션이 아니라 도착 경로 자체에 날짜를 넣습니다.

from datetime import date

target_path = f"/work/incoming/{date.today():%Y/%m/%d}"

단계 연결

수집이 끝나면 결과 전달이 자동으로 시작되게 하기

각 구간을 자동화로 만들고 같은 flowId로 묶습니다. 다음 단계의 triggerAutomation.value에 앞 단계의 automationId를 넣으면 서버가 이어서 실행합니다.

import uuid

flow_id = str(uuid.uuid4())


def build_step(name, source, source_path, target, target_path,
               step, trigger_id=None, options=None, webhook=None):
    schedule = {
        "type": "none",
        "startDateType": "now",
        "startDate": now_iso(),
        "timezone": "Asia/Seoul",
    }

    if trigger_id:
        schedule["triggerAutomation"] = {"value": trigger_id}

    body = {
        "name": name,
        "flowName": name,
        "flowId": flow_id,
        "transferType": "normal",
        "timezone": "Asia/Seoul",
        "step": step,
        "isUpcoming": False,
        "details": [
            {
                "senderId": source,
                "receiverId": target,
                "sourceItem": [
                    {
                        "hash": encode_path(source, source_path),
                        "filePath": source_path,
                        "isDir": True,
                    }
                ],
                "targetPath": encode_path(target, target_path),
                "step": step,
                "transferOptions": {
                    "noSchedule": False,
                    "target-action": "numbering",
                    "send-fileoption": {},
                    **(options or {}),
                },
            }
        ],
        "schedules": [schedule],
    }

    if webhook:
        body["processors"] = [{
            "category": "run",
            "type": "http",
            "config": {"url": webhook, "method": "POST"},
        }]

    return body


collect_id = api("POST", "/api/automations", build_step(
    "collect", "device-cam-01", "/media/raw",
    PROCESS_ID, "/work/incoming", step=1,
    options=FILTER, webhook=ENCODE_HOOK))["automationId"]

archive_id = api("POST", "/api/automations", build_step(
    "archive", PROCESS_ID, "/work/output",
    ARCHIVE_ID, "/archive", step=2,
    trigger_id=collect_id, options=TARGET_OPTIONS))["automationId"]

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

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

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

한 단계의 소스와 대상은 서로 다른 장비여야 합니다. 같은 장비를 지정하면 400이 반환됩니다.

프로세서의 호출 시점은 전송 완료 후입니다. config.events로 어느 이벤트에 반응할지 지정하며, 실패까지 받으려면 받는 쪽에서 상태를 확인합니다.

변경분 전송

마지막 실행 이후 바뀐 파일만 보내기

transfer = api("POST", "/api/transfers/manual", {
    "sourceDevice": "device-cam-01",
    "targetDevice": PROCESS_ID,
    "targetPath": "/work/incoming",
    "sourcePaths": ["/media/raw"],
    "sendAllFolder": True,
    "incremental": True,
    "transferOptions": {"target-action": "overwrite"},
})

기본값은 꺼짐입니다. 변경분 계산은 장비의 에이전트가 수행하며, 마지막 실행 이후 추가되거나 수정된 파일만 전송합니다.

증분 전송에서 도착 정책은 반드시 overwrite로 둡니다. 번호를 붙이는 정책이면 수정된 파일이 새 이름으로 쌓여 처리 장비가 예전 파일을 계속 읽게 됩니다.

결과 확인과 재전송

단계별 회차를 확인하고 실패한 파일만 재전송하기

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 step, automation_id in enumerate([collect_id, archive_id], start=1):
    runs = api("GET", f"/api/automations/{automation_id}/executions") or []
    latest = runs[0] if runs else {}

    if latest.get("status") not in (None, STATUS_COMPLETE):
        print(f"step {step} failed, retried {retry_failed(latest['monitorId'])} files")
        break

실행 이력은 최신 회차가 배열 앞에 오며 페이징 없이 전체가 반환됩니다.

확인 항목확인 내용
선별 조건확장자·크기·이름 필터
수집장비별 도착 경로 분리
저장 위치savepathoptionPath
흐름단계별 회차와 상태
재전송실패한 파일과 처리 결과