Skip to content

스크립트 개발 / 커넥터 구독

일부 커넥터는 메시지 구독을 지원하며, DataFlux Func는 통일된 방식으로 구독할 수 있도록 제공합니다.

이전 버전에서는 '데이터 소스'라고 불렸으며, 현재 버전에서는 '커넥터'로 변경되었습니다

1. 서문

DataFlux Func의 스크립트 실행 메커니즘 때문에 함수는 시작된 후 반드시 종료되어야 하며, 함수가 무한히 실행되는 것은 허용되지 않습니다.

따라서 구독과 같은 장기 상주 처리는 DataFlux Func에서 소비자를 직접 작성하여 장기간 실행하는 방식으로는 지원되지 않습니다.

대신 커넥터에서 구독 Topic을 지정하고, DataFlux Func 메인 프로그램이 구독 메시지를 일괄적으로 담당하도록 해야 합니다.

DataFlux Func 메인 프로그램이 메시지를 수신하면, 메시지를 지정된 메시지 처리 함수로 전달하여 처리함으로써 구독 메시지 처리를 완료합니다.

2. 구독을 지원하는 커넥터

최신 버전의 DataFlux Func는 다음 커넥터의 구독을 지원합니다

3. 작업 절차

Redis 메시지를 구독하고 처리하는 구체적인 작업 절차는 다음과 같습니다.

3.1 메시지 처리 함수 작성

메시지 처리 함수는 다음과 같은 고정된 함수 형태를 가집니다.

Python
1
2
3
4
@DFF.API('Message Handler')
def message_handler(topic, message):
    print('topic', topic)     # 주제
    print('message', message) # 내용

스크립트 작성이 끝나면 저장하고 게시합니다.

message-handler-func.png

3.2 커넥터 생성 및 구성

기본 구성 외에도 구독할 Topic과 해당 처리 함수를 추가로 입력해야 합니다.

처리 함수는 위에서 작성한 message_handler(...) 함수입니다.

sub-config.png

3.3 메시지 게시 및 메시지 처리 확인

발행자가 다음과 같은 메시지를 발행하면, 메시지는 DataFlux Func 메인 프로그램을 통해 위의 message_handler(...) 함수로 전달되어 처리됩니다.

커넥터 구성 페이지의 해당 Topic 아래에 최신 소비 정보가 표시됩니다.

sub-info.png

클릭하면 더 자세한 작업 정보를 확인할 수 있습니다.

sub-info-detail.png

이로써 Redis에 게시된 메시지가 실제로 message_handler(...) 함수에 의해 수신되고 처리되었음을 확인할 수 있습니다.

4. 메시지 게시

구독 메시지를 지원하는 커넥터는 일반적으로 메시지 게시(publish)도 지원합니다.

메시지 게시에 대한 자세한 내용은 스크립트 개발 / 커넥터 객체 DFF.CONN의 해당 커넥터 객체 API를 참조하세요. 일반적으로 .publish(topic, message) 형식입니다.

5. 메시지 처리 함수의 작업 기록

최신 버전의 DataFlux Func에서는 메시지 처리 함수의 작업 기록을 직접 확인할 수 있습니다.

message-handler-task-info.png

구독하는 메시지 수가 매우 많으면 각 메시지의 처리 로그를 기록하는 것 자체가 성능 문제를 일으킬 수 있습니다. 배포 및 유지보수 / 시스템 지표 및 작업 기록 / 로컬 함수 작업 기록 비활성화를 참조하여 '로컬 함수 작업 기록'을 비활성화하면 MySQL 저장 부담을 줄일 수 있습니다

구버전 DataFlux Func에서는 메시지 처리 함수가 작업 기록 조회를 지원하지 않습니다. 처리 중 발생한 오류를 기록하려면 메시지 처리 함수에서 관련 정보를 DFF.CACHE에 작성하면 됩니다.

다음 코드 구현을 참조하세요.

Python
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
import arrow
import traceback

def message_handler_impl(topic, message):
    # 실제 메시지 처리 함수
    # 0으로 나누기 오류가 발생한다고 가정
    x = 100 / 0

@DFF.API('Message Handler')
def message_handler(topic, message):
    try:
        # 실제 메시지 처리 함수 호출
        message_handler_impl(topic, message)

    except Exception as e:
        # 현재 시간 가져오기
        now_str = arrow.now('Asia/Shanghai').format('YYYY-MM-DD HH:mm:ss')

        # 전체 오류 정보 추출
        error_stack = traceback.format_exc()

        # 오류 정보를 DFF.CACHE에 저장
        latest_error = '\n'.join([ '시간: ', now_str, '오류 정보: ', error_stack ])
        DFF.CACHE.set('latest_error', latest_error)

        # 오류는 다시 던져야 함
        raise

실행 과정에서 문제가 발생하면 '관리 / 함수 캐시 관리자'에서 구체적인 정보를 확인할 수 있습니다.

6. 단일 구독 및 다중 구독

버전 1.7.31에서 추가됨

모든 구독기는 DataFlux Func의 server 서비스에서 실행됩니다. 구독이 메시지를 수신하면 구성에 따라 함수 실행 작업을 생성하여 worker로 보내 실행합니다.

구독기는 2가지 구독 방식을 지원합니다:

  1. 단일 구독: server 복제본을 몇 개 실행하든 구독기는 항상 그중 하나의 server에서만 실행됩니다
  2. 다중 구독: 여러 개의 server 복제본을 실행하면 구독기가 각 server에서 모두 실행됩니다

Redis 구독기는 공유 구독 처리를 지원하지 않으므로 다중 구독은 메시지 중복 소비만 발생시킵니다. 따라서 현재는 단일 구독만 지원합니다.

MQTT, Kafka 구독기는 단일 구독 또는 다중 구독을 선택할 수 있습니다.

MQTT를 다중 구독으로 사용할 때 Topic은 $share/group_name/topic_name 방식으로 공유 구독해야 합니다. 그렇지 않으면 동일한 메시지를 중복으로 수신하게 됩니다.

DataFlux Func가 포함하는 서비스와 여러 server 복제본 실행에 대한 자세한 내용은 배포 및 유지보수 / 아키텍처, 확장 및 리소스 제한을 참조하세요.

7. 구독 처리 속도 향상

버전 1.7.31에서 추가됨

DataFlux Func의 경우 구독기는 server 서비스에서 실행되고 함수는 worker 서비스에서 실행되므로 구독 처리 속도는 두 부분으로 구성됩니다: 구독 메시지 수신 속도와 구독 함수 실행 속도입니다.

구독 메시지 수는 많지 않지만 각 메시지 처리가 복잡하고 시간이 오래 걸리는 경우, 메시지가 적시에 처리되도록 worker 수를 늘려야 합니다.

구독 메시지 수가 매우 많은 경우, serverworker 수를 동시에 늘리고 구독기를 '다중 구독' 모드로 설정해야 합니다.

시스템 확장 방법에 대한 자세한 내용은 배포 및 유지보수 / 아키텍처, 확장 및 리소스 제한을 참조하세요.

8. 구독 제한

버전 1.7.31에서 추가됨

구독 메시지 수가 매우 많을 때, 서버 측의 게시-구독 방식이 다르기 때문에 그에 따라 서로 다른 제한이 존재합니다.

Redis / MQTT

Redis, MQTT 구독기는 구독 측에서 메시지 수신 속도를 제어할 수 없으므로 수신된 메시지는 내부 버퍼 풀로 들어갑니다.

버퍼 풀의 기본 크기는 5000이며, 최대 5000개의 구독 메시지가 처리를 기다리며 머무를 수 있습니다. 그 사이에 더 많은 메시지가 수신되면 초과된 메시지는 삭제됩니다.

위 문제를 해결하려면 worker 수를 늘려 구독 메시지 처리 속도를 높일 수 있습니다.

Kafka

Kafka 구독기는 구독 측에서 메시지 수신 속도를 제어할 수 있으므로 버퍼 풀이 존재하지 않으며 메시지도 삭제되지 않습니다.

메시지 처리 속도가 메시지 게시 속도를 따라가지 못할 때, worker 수를 늘려 구독 메시지 처리 속도를 높일 수 있습니다.

9. 알려진 문제

현재 알려진 문제는 다음과 같습니다.

MQTT 끊김 및 메시지 수신 불가

MQTT Broker가 EMQX에 연결할 때 끊김이나 메시지 수신 불가 문제가 발생할 수 있습니다.

이 문제는 초기 커뮤니티 단일 서버 버전의 EMQX에서 발생한 적이 있으며, 다른 버전의 EMQX는 아직 테스트되지 않았습니다.

그러나 mosquitto에 연결할 때는 이 문제가 나타나지 않았습니다.

Kafka 시작 후 즉시 소비할 수 없음

Kafka 커넥터는 구독 후 처음 몇 분 동안 일시 중지 상태에 있을 수 있으며, 이후에야 정상적으로 메시지를 소비하게 됩니다.

원인은 Kafka의 Rebalance 메커니즘과 관련이 있을 수 있습니다.