Skip to main content

jetstream-api

Python으로 작성된 ILink 클라이언트 API입니다.

시스템 요구 사항

시스템에 Python 3.9 이상, pip가 설치되어 있어야 합니다. 필수 의존성은 없습니다(표준 라이브러리만 씁니다). lz4 배치 압축을 쓰려면 pip install "jetstream-api[lz4]", zstd 배치 압축을 쓰려면 pip install "jetstream-api[zstd]" 로 설치합니다(zstd 는 Python 3.14 이상이면 표준 라이브러리로 됩니다). 둘 다 쓰려면 pip install "jetstream-api[lz4,zstd]" 입니다.

설치 방법

pip install jetstream-api

0.17.0 변경 - 클러스터 API 를 컨트롤러 정족수 모델로 (자바 v2.13.0 과 같음, 엔진은 해당 판 이상)

판사 · MAIN · SUB 세 정의를 각각 만들던 모델이 없어졌습니다. 이제 클러스터 하나에 정의 하나이고, 그 안에 노드 전부를 적습니다. 역할은 정의의 종류가 아니라 노드의 성질입니다 - 큐 관리자가 있으면 브로커, coordinator 를 세우면 컨트롤러이고, 둘은 겹칠 수 있습니다.

c = ILClusterProperty("C1")
c.nodes = [
    ILClusterNode("QM1", "10.10.1.12", 12345, 9000, node_id=1),
    ILClusterNode("QM2", "10.10.1.13", 12345, 9000, node_id=2),
    ILClusterNode("QM3", "10.10.1.14", 12345, 9000, node_id=3),
]
for n in c.nodes:
    n.coordinator = True          # ilsc 의 COORDS(1,2,3)
svc.createCluster(c)
  • 컨트롤러 포트가 노드마다 하나입니다(ILClusterNode.ctrlPort). 클러스터 한 개에 하나였던 ilccPort 는 없어졌습니다.
  • 노드 번호(nodeId)는 자리 이름이지 위치가 아닙니다. 노드를 지우면 번호에 구멍이 남고 메우지 않습니다 - nodes[n-1] 로 찾지 말고 getNode(n) 을 쓰십시오.
  • 큐 관리자 없이 투표만 하는 노드를 만들 수 있습니다: ILClusterNode(None, "10.10.1.15", 0, 9000, node_id=4).
  • getBrokerNodes() / getControllerNodes() / getNode(n) / ILClusterNode.getRoles() 가 새로 있습니다.
  • coordinator 는 참/거짓이 아니라 세 상태입니다 - True(맞다) / False(아니다) / None(모름). is None 으로 먼저 가르십시오. isCoordinator() 도 세 값을 그대로 줍니다.
  • replicas / minSync 는 비워 두는 것이 기본이고, 비운 채로 왕복합니다. 비우면 엔진이 브로커 역할 노드 전부를 복제본 수로 치므로 노드를 더하면 따라 늘어납니다. 숫자를 명시하면 그 진술이 고정되어 노드를 더할 때 같이 고쳐 보내야 합니다. 0 은 정상값이고 "안 정함" 은 None 뿐입니다 - if not p.replicas 로 판정하지 마십시오.
  • ☠ minSync 는 Kafka 의 min.insync.replicas 와 이름만 같고 계약이 다릅니다. ISR 이 이 아래로 떨어져도 쓰기를 거절하지 않습니다 - 가용성을 택하고 내구성을 내립니다.
  • setClusterProperty 는 읽어서 고쳐 되보내십시오. 정의에는 클라이언트가 읽기만 하는 자리(clusterId · signature · (defTerm, defVersion))가 있고, getClusterProperty 로 받은 객체가 그 바이트를 그대로 들고 있습니다. 새로 만든 객체로 ALTER 하면 그 자리가 빈칸으로 나갑니다.
  • createCluster 는 더 이상 클라이언트에서 규칙을 미리 검사하지 않습니다. 규칙은 엔진이 정본이고, 클라이언트에 한 벌 더 적으면 두 벌이 갈립니다(옛 TYPE 사전검사가 그 경로로 생겼습니다).
  • 와이어 레이아웃이 바뀌었습니다 - 고정부 512 → 768, 노드 256 → 384(ILC.MIMQ_VO_CLUSTER_LEN / MIMQ_VO_CLUSTER_NODE_LEN). 블롭 앞머리에 판 표시 "V2" 가 붙고, 표시가 없는 옛 블롭은 거절합니다(길이로는 가릴 수 없습니다 - 옛 512+5+256 과 새 768+5+0 이 둘 다 773 입니다).
  • 없어진 이름: ILClusterProperty.judge() · ILClusterProperty.coordinator()(클래스 메서드) · clusterType / getClusterType / setClusterType · ilccPort / getIlccPort / setIlccPort(클러스터와 노드 양쪽) · authorityHolder / getAuthorityHolder / getAuthorityHolderType · isJudge · ILClusterNode.get_ilcc_address / getIlccAddress · ILClusterProperty(name, ilcc_port, is_auto_start, cluster_type) 의 ilcc_port · cluster_type 인자. 없어진 이름에 대입하면 무엇으로 바뀌었는지까지 알려 주는 AttributeError 가 납니다.

0.16.0 변경 - 세그먼트 이름을 토픽과 같은 형식으로, 오타 대입 거부 (자바 v2.12.0 과 같음)

  • 이름 규칙이 정해졌습니다 - 개수는 맨 복수형, 목록은 ...List 접미. 토픽이 partitions 를 쓰는 형식 그대로입니다. 0.15.0 이 하루 쓴 이름 둘을 여기로 옮깁니다.

    자리 0.15.0 0.16.0
    세그먼트 개수 (ILQueueProperty) numberOfSegment segments
    세그먼트 목록 (ILQueueRealTimeInfo) segments segmentInfoList
    • ★ 0.15.0 의 두 이름은 남기지 않고 지웠습니다. 남기면 개수에 이름이 셋이 되고, 더 나쁘게는 segments 가 클래스마다 다른 뜻(개수 vs 목록)이 됩니다. 0.15.0 은 하루짜리였고, 아래 가드 덕분에 지운 것이 조용하지 않습니다.
    • numberOfPartition 과 partitions 는 그대로 있습니다. 예전부터 쓰던 이름이라 영구 별칭으로 남습니다.
    • 큐에 진짜 분산 파티션이 생기면 그 수는 토픽과 같은 이름(partitions) 으로 따로 들어옵니다. 지금 있는 이름을 그쪽으로 돌리지 않습니다 - 재사용은 제거보다 나쁩니다. 제거는 그 자리에서 터지지만 재사용은 조용히 다른 값을 줍니다.
  • 속성 객체가 모르는 필드 이름 대입을 거부합니다. 지금까지는 prop.numberOfSegmnet = 4 같은 오타가 예외도 없이 쓰레기 속성을 만들고 실제 값은 그대로였습니다 - 읽기는 원래 AttributeError 였는데 쓰기만 침묵해서 옛 값이 그대로 서버로 나갔습니다.

    • 이름이 바뀐 필드는 무엇으로 바뀌었는지까지 알려 줍니다: AttributeError: ILQueueProperty has no attribute numberOfSegment - 0.16.0 에서 segments 로 바뀌었다
    • 정말 임의 속성을 얹어야 하면 object.__setattr__(obj, name, value) 로 우회할 수 있습니다. 그 값은 와이어에 실리지 않습니다.
    • 생성자가 도는 동안은 열려 있습니다(필드를 만들어야 하므로). 상속해 쓰는 응용도 super().__init__() 뒤에 자기 필드를 만들 수 있습니다.
  • 와이어는 달라지지 않습니다. 0.15.0 과 바이트가 같습니다(큐 속성 2200 · 구엔진 4096 · 실시간정보). 이 판도 엔진 판을 가리지 않습니다.

0.15.0 변경 - 큐의 "파티션" 이 "세그먼트" 가 됩니다 (자바 v2.11.0 과 같음)

  • ILQueueProperty.numberOfPartition 이 numberOfSegment 로 바뀝니다. 값도 자리도 그대로이고 이름만 바뀝니다.

    • 옛 이름은 그대로 돕니다. numberOfPartition 은 같은 값을 가리키는 별칭으로 남습니다 - 읽기도 쓰기도 되고, getNumberOfPartition() / setNumberOfPartition() 도 전처럼 동작합니다. 고치지 않아도 되는 변경입니다.
    • __str__ 라벨이 NUM-PARTITION 에서 NUM-SEGMENT 로 바뀝니다. 이 출력을 정규식으로 읽는 곳이 있으면 그곳만 봐 주세요.
  • 왜 바꾸는지: 같은 "파티션" 이라는 말이 토픽과 큐에서 정반대 를 뜻하고 있었습니다.

    말 뜻
    파티션 분산 처리 단위 (토픽. 리더 · 팔로워가 붙는다)
    세그먼트 한 노드 안의 저장 칸 (큐. 파일 1~N, 메모리, 공유메모리)

    토픽의 PARTITIONS 와 큐의 PARTITION 이 한 글자 차이로 서로 다른 개념이었습니다. 엔진이 큐에 진짜 분산 파티션을 들이기 전에 옛 이름을 비워 두는 것입니다.

  • 와이어는 달라지지 않습니다. 큐 속성 직렬화는 위치 기반 고정폭이라 속성 이름이 실리지 않습니다. 어느 이름으로 넣든 나가는 2200바이트(구엔진용 4096 레이아웃도)가 같습니다. 그래서 이 판은 엔진 판을 가리지 않습니다.

  • 엔진 쪽 이름도 같이 바뀝니다(ilsc PARTITION(n) -> SEGMENT(n), VIEW QUEUE 머리글, 컨피그 XML Storage/@NumberOfPartition -> @NumberOfSegment). 엔진도 옛 이름을 계속 받습니다.

    • ★ ilsc 에서 SEGMENT 를 쓰려면 엔진도 같이 올려야 합니다. 클라이언트만 올린 환경에서는 엔진이 그 키워드를 모릅니다 - 그때도 큐는 정상 동작하고, ilsc 에서 PARTITION 을 쓰면 됩니다.
  • 큐 실시간 정보의 "파티션" 목록도 같이 바뀝니다. getQueueRealTimeInfo() 가 돌려주는 ILQueueRealTimeInfo 의 partitions 가 segments 가 되고, 항목 클래스 PartitionItem 은 SegmentItem 이 됩니다.

    • 여기도 옛 이름이 그대로 돕니다. partitions 는 같은 목록을 가리키는 별칭이고(getPartitions() / setPartitions() 포함), PartitionItem 은 SegmentItem 과 같은 클래스 객체라 isinstance 도 전처럼 통합니다.
    • ★ 토픽의 getPartitions() 와는 뜻이 다릅니다. 지금까지 두 클래스에서 같은 이름이 정반대를 뜻하고 있었습니다 - 토픽 쪽은 분산 처리 단위, 큐 실시간 정보 쪽은 한 노드 안의 저장 칸입니다. 이번 개명의 목적이 바로 이것입니다.
    • __str__ 머리글이 Partitions --- 에서 Segments --- 로 바뀝니다(엔진 VIEW QUEUE 가 SEGMENTS( 로 나가는 것과 맞춥니다).
    • ILC.MIMQ_VO_SEGMENT_ITEM_LEN 이 생기고 ILC.MIMQ_VO_PARTITION_ITEM_LEN 은 같은 값(128)의 별칭으로 남습니다.
  • 자바 클라이언트는 클래스 이름을 바꾸지 않습니다(PartitionItem 유지, getSegments() 추가). 자바는 제네릭이 불변이라 ArrayList<PartitionItem> 을 별칭으로 덮을 수 없고, 사본으로 옛 접근자를 만들면 목록 수정이 조용히 사라지기 때문입니다. 파이썬은 클래스 별칭이 완전해서 이런 제약이 없습니다.

  • 없어진 이름은 하나도 없습니다. 이 판은 추가와 별칭뿐입니다.

  • 디스크 파일 이름은 그대로입니다 (QUEUE.P001). 저장 레이아웃 자체라 바꾸지 않습니다.

0.14.0 변경 - 소켓 keepalive 탐침 타이머, NIO 커넥터 수신 기한 정합 (자바 v2.10.0 과 같음)

  • 연결 소켓에 keepalive 탐침 타이머를 겁니다(유휴 60초 + 10초 간격 x 6회). 엔진이 쓰는 값과 같습니다. 지금까지는 SO_KEEPALIVE 만 켜고 타이머는 운영체제 기본에 맡겼습니다(리눅스 약 2시간 11분).
    • 얻는 것은 주로 NAT · 방화벽 유휴 만료 방지입니다. 60초마다 탐침이 흐르므로 중간 장비가 연결 상태를 지우지 않습니다. 조용한 발행자가 밤새 붙어 있는 배치에서 제일 자주 걸립니다.
    • 상대가 사라진 연결을 커널이 약 2분에 끊어 주는 것은 덤입니다. 다만 라이브러리가 그것을 감지해 스스로 다시 붙지는 않습니다 - 실패는 다음 요청에서 드러나고, 그건 이미 세션 타임아웃(기본 30초)이 덮습니다.
    • 응용 요청이 아니라 커널 탐침입니다. 엔진의 수신 타임아웃(SOCKTIMEOUTRCV)은 요청을 읽는 시계라 탐침으로는 갱신되지 않습니다. 조용한 세션은 그 기한을 향해 그대로 늙습니다.
    • 타이머를 걸지 못해도 연결은 그대로 씁니다. 그때는 운영체제 기본 타이머로 동작합니다. 결과는 ilink.connectivity.keepalive_timer_state 로 볼 수 있습니다(진단용).
  • NIO 커넥터(TCPConnectorNIO)의 수신 기한을 나머지와 같은 규칙으로 맞췄습니다. 기한이 프레임 전체가 아니라 소켓 읽기 하나마다이고, 0 이하면 끝없이 기다립니다. 0.13.0 까지는 헤더와 본문이 한 기한을 나눠 썼습니다.
    • 송신 쪽은 전처럼 기한 하나를 씁니다. 조금씩만 나아가는 송신이 끝나지 않던 것을 0.7.4 에서 막은 것입니다.
    • 이 커넥터는 set_conn 으로 직접 갈아 끼울 때만 쓰입니다. 기본 경로는 달라진 것이 없습니다.

keepalive 타이머가 실제로 걸리는 환경

자바 클라이언트는 걸 수 있는 자리가 좁습니다. TCP_KEEPIDLE 계열이 jdk.net.ExtendedSocketOptions(자바 11 이상)에만 있고, 윈도우는 그 판에서도 받지 않기 때문입니다.

자바 8 자바 11 이상
윈도우 걸리지 않음 걸리지 않음
리눅스 · macOS 걸리지 않음 걸림

파이썬 클라이언트는 전 플랫폼에서 걸립니다. 윈도우는 탐침 횟수가 운영체제 고정값(10회)이라 약 160초, 그 밖은 60초 + 10초 x 6 = 약 2분입니다.

0.13.0 변경 - 세션 타임아웃 기본값 30초, acks 값 검사, 수신 기한 수정 (자바 v2.9.0 과 같음)

  • 세션 타임아웃 기본값이 경로마다 달라집니다. 리더가 얼어붙어 응답이 오지 않으면 0.12.0 까지는 최대 10분을 매달렸습니다.
    • 데이터 경로(ILQmgr / 큐 / 토픽): 600000 -> 30000. 접속 대기 시간도 같은 값이라 페일오버 재접속이 빨라집니다.
    • 관리 경로(ILAdminService): 600000 유지. 엔진이 큐 관리자 정지 · 삭제 응답을 실제로 내려갈 때까지 붙들고, 토픽을 닫으며 파티션마다 fsync 하느라 수 분이 걸릴 수 있습니다.
    • 큐 get(wait) 은 전처럼 알아서 늘어납니다(대기 시간의 2배). 기본 대기 60초면 수신 상한이 120초라 기본 동작은 그대로입니다. 무한 대기(set_wait_msg_timeout_ms(-1))도 이제 같이 늘어납니다. 무한 대기는 엔진에 60초를 주고 루프를 도는 것이라 엔진이 쥐는 시간은 유한 대기와 같은데 전에는 예외로 두고 있었습니다. 기본값이 600000 일 때는 드러나지 않던 자리입니다.
    • 전처럼 오래 기다리려면 qmgr.set_session_timeout(600000) 을 부르면 됩니다.
  • spokeSwitchTo(key, timeoutSec) 는 엔진이 그 시간만큼 응답을 붙들기에, 그 호출 동안만 수신 상한을 timeoutSec 의 2배로 올리고 되돌립니다.
  • 수신 기한이 소켓 읽기 하나마다 걸립니다. 프레임 전체의 기한이 아닙니다 - 응답이 조각으로 나뉘어 와도 다음 조각이 그 시간 안에 오면 계속 받습니다(자바와 같은 규칙). 느린 회선의 큰 메시지가 파이썬에서만 끊기던 것이 사라집니다.
    • 기한이 지나면 MIMQE_WAIT_TIMEOUT 입니다. 전에는 MIMQE_INTERNAL_ERROR("unhandled exception")로 나왔습니다.
    • 세션 타임아웃 0 은 끝없이 기다립니다(자바와 같습니다). 전에는 곧바로 기한 초과였습니다.
  • ILProducerConfig.acks 는 0 과 1 만 받습니다. 그 밖의 값은 ValueError 입니다(자바 IllegalArgumentException 과 짝). 0.12.0 까지는 0 이 아닌 값을 모두 1 로 보정해 acks=2 같은 오타가 조용히 통과했습니다.
  • 문서: 클러스터에서 acks(1) 은 팔로워 확인(HW)까지라는 것(= Kafka acks=all), 파티션이 많은 토픽에서는 그 대기가 봉투 안 파티션마다 겹쳐 requestTimeoutMs 가 먼저 끝날 수 있다는 것, removeSession 이 반쯤 열린 연결을 푸는 수단이라는 것을 적었습니다.

0.12.0 변경 - 큐 전용 요청 코드(엔진 7.0.1.3431 이상) 자동 사용 (자바 v2.8.0 과 같음)

  • 앱 코드는 바꿀 것이 없습니다. 엔진 판 줄이 7.x 이고 리비전이 3431 이상이면 클라이언트가 스스로 새 요청 코드를 씁니다. 그 밖의 엔진(6.3.x, 그보다 앞선 7.x, 5.x)에는 0.11.0 과 같은 요청을 보냅니다.
  • commit_and_get 은 커밋 요청과 get 요청 대신 전용 요청 하나(COMMIT_AND_GET, 3402)를 보내고 응답도 하나만 받습니다. 엔진이 커밋한 뒤 곧바로 get 을 돌립니다. 결과와 예외는 0.11.0 과 같습니다.
    • 커밋이 거절되면 get 은 돌지 않았고, 롤백을 보내 커밋하려던 작업과 꺼낸 메시지를 큐에 되돌린 뒤 ILOperationException 을 올립니다.
    • 응답을 받기 전에 연결이 끊기면 커밋 결과를 알 수 없습니다(재접속 / 재시도 없음). 커밋 결과가 get 결과와 한 응답으로 오므로 결과를 알 수 없는 구간이 get 대기 시간까지 늘어납니다.
  • commit() 은 쥐고 있던 트랜잭션 put 이 하나뿐이면 put 과 커밋 대신 전용 요청 하나(PUT_AND_COMMIT, 3403)로 보냅니다. 돌려주는 값과 예외는 같습니다. put 이 둘 이상이거나 puts, 64KiB 를 넘는 put, 쥔 것이 없을 때는 0.11.0 과 같습니다.

0.11.0 변경 - 큐 단건 put / get 가속, commit_and_get 추가

  • 새 API ILQueue.commit_and_get(matchId=None, option=MatchOption.NONE) (camelCase 별칭 commitAndGet): 이 세션의 미확정 작업을 커밋하고 이어서 다음 메시지를 꺼냅니다. qmgr.commit() 뒤에 q.get() 을 부른 것과 결과가 같고, 커밋 요청과 get 요청을 응답을 기다리지 않고 연달아 보내므로 1 왕복입니다(따로 부르면 2 왕복). get → commit 을 반복하는 소비 루프는 반드시 commit_and_get 을 쓴다. 확정 시점은 그대로(처리를 마치고 부르는 순간)입니다.

    q = qmgr.access_queue("APP.REQUEST.Q")      # 트랜잭션 세션
    try:
        msg = q.get()
        while True:
            process(msg)
            msg = q.commit_and_get()            # 처리한 것을 확정하고 다음 것을 꺼낸다 - 1 왕복
    except ILNoMsgException:
        pass                                    # 마지막 커밋은 이미 끝났다
    
    • 대기 시간 안에 메시지가 없으면 ILNoMsgException 입니다 - 이때도 커밋은 이미 성공했습니다. 커밋한 메시지 수는 돌려주지 않습니다(필요하면 commit() 을 따로 부릅니다).
    • 커밋이 실패하면 commit() 과 같은 ILOperationException 입니다. 커밋 뒤에 보낸 get 이 꺼낸 메시지는 롤백으로 큐에 되돌린 뒤 올립니다. 커밋은 됐는데 get 을 엔진이 거절하면 get() 과 같은 예외입니다.
    • 커밋 요청을 보낸 뒤 커밋 응답 전에 연결이 끊기면 커밋 결과를 알 수 없으므로 자동 재접속 / 재시도 없이 ILSessionException 입니다(commit() 과 같습니다). 커밋이 성공한 뒤의 get 장애는 get() 처럼 자동 재접속 후 1회 재시도합니다.
    • 자동 커밋 세션이면 커밋 요청 없이 get() 과 똑같습니다. 무한 대기(set_wait_msg_timeout_ms(-1))도 get() 과 같습니다.
  • 와이어 형식은 그대로입니다. 5.x / 6.x / 7.x 엔진 모두에서 씁니다. 연달아 보내기는 한 세션의 요청을 받은 순서대로 처리하고 그 순서로 답하는 것이 확인된 엔진(7.x, 그리고 6.3.x 중 6.3.3 부터)에서만 하고, 그 밖의 엔진(5.x 등, 판을 읽지 못한 경우 포함)에서는 커밋 응답을 받은 뒤 get 을 보냅니다(2 왕복, 결과는 같습니다).

  • 내부 가속 - 공개 API 변화는 commit_and_get 추가뿐입니다. 보내는 바이트는 0.10.0 과 같습니다(옛 판과 바이트 대조 시험).

    • put 프레임을 한 번에 조립합니다(본문 복사 한 번). JMS 헤더 / 응답 칸은 한 번에 풉니다.
    • 받기: 연결마다 수신 버퍼에 recv_into 로 받고, 소켓 timeout 은 값이 바뀔 때만 줍니다. 큰 get 응답은 본문을 한 번만 복사합니다(getData() 는 여전히 호출자 소유의 bytes).
    • COMMIT / ROLLBACK 요청 프레임과 큐 핸들의 GET 요청 프레임을 짜 두고 씁니다(대기 시간 / 선택자가 바뀌면 다시 짭니다).
    • 트랜잭션 세션의 put(프레임 64KiB 이하)은 곧바로 보내지 않고 이 연결의 다음 요청(보통 commit())과 한 번에 보냅니다. 요청 순서 처리가 확인된 엔진(7.x, 6.3.3 부터)에서만 그렇게 하고, 그 밖의 엔진에는 0.10.0 과 같이 곧바로 보냅니다. 그래서 연결 오류가 put 이 아니라 그 다음 요청에서 드러납니다. rollback() 은 아직 보내지 않은 put 을 ROLLBACK 과 한 번에 보냅니다(돌려주는 값은 0.10.0 과 같습니다). 연결이 닫힌 것을 이미 알면 예전처럼 put 에서 바로 실패합니다. NIO 커넥터 (set_conn(TCPConnectorNIO()))는 묶어 보내지 않습니다.
    • client_name 을 비운 connect() 의 호출 모듈 이름 판별이 가벼워졌습니다(결과는 같습니다).
  • put 이 채우는 msgId: 여전히 uuid4 모양(36자, 소문자 16진)이지만 random 모듈의 128비트로 만듭니다. 유일 식별자이지 보안 토큰이 아닙니다 - 예측할 수 없어야 하면 msgId 를 직접 지정하세요. fork 한 자식 프로세스는 random 이 다시 시드되므로 부모와 같은 값을 내지 않습니다.

  • 사용 안내: 큐 핸들은 한 번 얻어 계속 씁니다(메시지마다 access_queue() 를 부르면 존재 확인 왕복이 한 번씩 더 듭니다 - 존재를 이미 알면 access_queue(name, exist_check=False)). ILQmgr 는 스레드마다 따로 둡니다.

0.10.0 변경 - 토픽 레코드 v5: 길이 칸 varint (자바 v2.6.0 과 같음)

  • 레코드 v5: producer 가 파티션 배치를 magic 5 로 씁니다. 레코드의 길이 칸 셋(key / 속성 / 값)이 4바이트 고정에서 unsigned LEB128(길이 + 1, 1~5바이트)로 바뀌었습니다 - 126바이트까지 1바이트, 16,382바이트까지 2바이트입니다. 배치 헤더 52바이트, CRC, 압축 코덱과 압축 범위, 오프셋, 멱등 seq 는 그대로입니다.
  • 레코드당 저장 바이트 감소: 값 100B + uuid 키 36B + 속성 없음 = 148 → 139B(-6.1%), 값 1KB = 1,072 → 1,064B(-0.75%). 압축을 켜면 차이는 거의 없습니다.
  • 소비는 배치마다 magic 을 봅니다: 4(i32 길이 칸)와 5 를 모두 풀고, 한 응답에 섞여 있어도 됩니다.
  • 깨진 v5 레코드: 길이 칸이 5바이트를 넘거나, 적힌 값이 2^31 이상이거나, 잘렸거나, 값이 없으면(적힌 값 0) 그 배치는 0.8.0 의 "풀지 못한 배치" 규칙대로 그 자리에서 멈춥니다(MIMQC_FATAL_ERROR / MIMQE_INTERNAL_ERROR). 사유는 invalid v5 varint (...) / invalid v5 record value (absent ...) 입니다.
  • 공개 API 변화 없음: 설정 / 메서드 / 반환 모양이 그대로입니다. 크기 검사(maxRequestSize / maxFrameLength / bufferMemory / batchSize)는 v5 부호화 길이로 셉니다 - 같은 설정에 레코드가 조금 더 들어갑니다.

0.9.0 변경 - 압축 코덱 lz4 / zstd, gzip 제거 (자바 v2.5.0 과 같음)

  • 코덱은 none / lz4 / zstd 입니다: gzip 을 빼고 lz4 를 넣었습니다(사용자 결정). 클라이언트 CPU·지연이 중요하면 lz4, 디스크·용량이 중요하면 zstd 를 고르세요. 기본값은 여전히 압축하지 않음(none)입니다.
  • 공개 API 제거: compressionGzipLevel() / compression_gzip_level() / get_compression_gzip_level() / getCompressionGzipLevel(). compressionType("gzip") 은 ValueError (compressionType must be none / lz4 / zstd (got 'gzip' - gzip and snappy are not supported))입니다.
  • 추가: compressionType("lz4"), compressionLz4Level(n) (compression_lz4_level) - 117, 기본 9, 게터 get_compression_lz4_level() / getCompressionLz4Level(). 뜻은 Kafka compression.lz4.level 과 같습니다 - 9(기본)는 빠른 압축기이고, 9 가 아닌 값은 그 레벨의 LZ4HC(훨씬 느리고 조금 더 줄입니다)입니다. 레벨 12 이상은 결과가 같습니다. 파이썬 lz4 모듈은 레벨 12 를 빠른 압축기와 같게 다룹니다(자바 lz4-java 는 1~2 도 HC 계열 - 라이브러리 차이).
  • lz4 는 선택 설치입니다: pip install "jetstream-api[lz4]"(lz4>=4.0). 모듈이 없으면 compressionType("lz4") producer 생성이 연결하기 전에 ILOperationException(MIMQE_NOT_SUPPORTED (compressionType lz4: no lz4 module - pip install "jetstream-api[lz4]"))으로 멈추고, lz4 배치를 받은 구독 read 는 원인 MIMQE_NOT_SUPPORTED (codec 3 lz4: no lz4 module - ...) 로 그 배치에서 멈춥니다(zstd 모듈 없음과 같은 MIMQC_USER_FAULT / MIMQE_NOT_SUPPORTED). 발행하는 코덱의 모듈을 구독하는 쪽에도 설치하세요.
  • lz4 와이어: 배치 attributes 코덱 3, 표준 LZ4 frame(lz4 명령과 같은 틀)입니다. producer 는 64KB 독립 블록 + 원문 크기로 만들고 체크섬은 넣지 않습니다(배치 CRC 가 압축 바이트를 덮습니다) - 자바 v2.5.0 도 같은 틀로 만듭니다. 소비는 이어 붙인 frame 을 모두 풀고, 잘림 / frame 뒤에 남은 바이트 / 푼 크기 128MiB 초과 / 깨진 frame 은 MIMQE_TOPIC_CODEC_ERROR (lz4: ...) 로 그 배치에서 멈춥니다(0.8.0 의 "풀지 못한 배치" 규칙 그대로).
  • 압축은 예전처럼 producer 의 I/O 스레드가 하고, 줄지 않는 배치는 압축하지 않고 보냅니다. lz4 모듈도 압축하는 동안 GIL 을 놓습니다.

0.8.0 변경 - 배치 압축 gzip / zstd (자바 v2.4.0 과 같음)

  • 엔진 요구: 압축 배치는 엔진 7.0.1.3419 이상(CMPR-1)에서만 받습니다. 그보다 옛 엔진은 압축 배치를 알아보지 못해 거절합니다. 기본값은 압축하지 않음(none)이라 설정하지 않으면 동작이 그대로입니다.
  • producer 설정 - 뜻과 기본값이 Kafka compression.type / compression.gzip.level / compression.zstd.level 과 같습니다:
    • compressionType("none" | "gzip" | "zstd") (compression_type) - 기본 "none". 소문자 그대로 씁니다(자바 / Kafka 와 같습니다). lz4 / snappy 는 지원하지 않습니다(ValueError).
    • compressionGzipLevel(n) (compression_gzip_level) - -1(zlib 기본 레벨 6, 기본값) 또는 1~9.
    • compressionZstdLevel(n) (compression_zstd_level) - -131072~22, 기본 3.
    • 게터 get_compression_type() / get_compression_gzip_level() / get_compression_zstd_level() (camelCase getCompressionType() / getCompressionGzipLevel() / getCompressionZstdLevel()).
    • 범위 밖 값은 빌더에서 곧바로 ValueError 입니다. 직접 대입(cfg.compressionType = "lz4")은 producer 생성 때 (validate()) 막습니다. 설정은 예전처럼 producer 생성 때 복사합니다.
  • zstd 는 선택 설치입니다: pip install "jetstream-api[zstd]"(zstandard>=0.22). 파이썬 3.14 이상이면 표준 라이브러리 compression.zstd 를 먼저 쓰므로 설치하지 않아도 됩니다. gzip 은 표준 라이브러리(zlib)만 씁니다 - 필수 의존성은 여전히 없습니다. zstd 모듈이 없으면 compressionType("zstd") producer 생성이 연결하기 전에 ILOperationException (MIMQE_NOT_SUPPORTED (compressionType zstd: no zstd module - pip install ...))으로 멈추고, zstd 배치를 받은 구독 read 는 ILException(원인 MIMQE_NOT_SUPPORTED (codec 4 zstd: no zstd module - pip install ...))입니다(자바와 같은 사유). 이 read 오류는 MIMQC_USER_FAULT / MIMQE_NOT_SUPPORTED 이고 그 배치에서 멈춥니다(아래 "풀지 못한 배치").
  • 발행: 파티션 배치를 보낼 때 레코드 열만 한 번 압축합니다(배치 헤더 52바이트는 그대로, recordCount 는 원래 건수, 배치 CRC 는 압축한 바이트가 대상). 재전송(연결 끊김, 순서 오류 재번호, PID 재발급)은 같은 압축 바이트에 헤더와 CRC 만 새로 씁니다. 압축해도 줄지 않는 배치(이미 압축된 자료, 무작위 바이트)는 압축하지 않고 보냅니다 - 한 토픽에 섞여도 됩니다. batchSize 와 레코드 크기 검사(maxRequestSize, maxFrameLength, 토픽 세그먼트, bufferMemory)는 압축 전 크기로 셉니다(Kafka 와 같습니다). 봉투에 배치를 채울 때는 압축한 뒤 길이로 세므로 봉투 하나에 배치가 더 실립니다.
  • 압축은 producer 의 I/O 스레드가 합니다: send() 는 예전처럼 배치 버퍼에 붙이기만 합니다. zlib / zstd 는 압축하는 동안 GIL 을 놓으므로 send() 를 부르는 스레드와 겹쳐 돕니다. 대신 producer 하나의 압축은 한 스레드라, 느린 코덱(gzip 레벨 6 은 텍스트 배치에서 대략 6070MB/s - 엔진 설계 문서의 측정)은 그 producer 의 발행 천장이 될 수 있습니다. zstd 나 낮은 gzip 레벨을 쓰거나 producer 를 늘리세요(한 토픽에 14개 안내는 그대로입니다).
  • 소비: 구독 read(배치 read 2042)와 패턴 구독(다중 토픽 read 2046)이 배치마다 코덱을 보고 풉니다. 한 토픽 / 한 응답에 무압축 / gzip / zstd 배치가 섞여도 됩니다. 오프셋과 커서 앞 레코드 빼기는 그대로입니다.
  • 풀지 못한 배치는 그 자리에서 멈춥니다(stop-in-place): 풀지 못하는 배치 - 모르는 코덱 MIMQE_TOPIC_UNSUPPORTED_CODEC (codec N), zstd 모듈 없음, 깨짐 / 잘림 / 뒤에 남은 바이트 / 푼 크기 128MiB 초과 MIMQE_TOPIC_CODEC_ERROR (gzip: ...) / (zstd: ...), 배치 CRC 불일치, 푼 레코드 수 / 길이가 헤더와 다름 - 를 만나도 응답 전체를 버리지 않습니다. 버리면 엔진 읽기 위치만 지나가 그 응답의 다른 배치까지 건너뛰고 다음 AUTO 커밋이 그 구간을 덮습니다(유실). Kafka 컨슈머처럼 그 자리에서 멈춥니다:
    • 그 배치(파티션 P, baseOffset B) 앞에 푼 레코드와 같은 응답의 다른 파티션 레코드는 여느 때처럼 줍니다(COMMIT_AUTO 커밋은 앱에 준 것만). 같은 응답의 P 뒤 배치는 버리고, P 는 구독 연결로 B 로 되감습니다(offset:B;partition:P seek - 커서가 B 보다 뒤면 커서로). 원인이 풀릴 때까지 P 는 서 있고 아무것도 건너뛰지 않습니다.
    • 오류는 read 에 ILException 으로 옵니다. 이번 호출에 줄 레코드가 있으면 그것을 돌려주고 다음 read 가 올리며, 없으면 곧바로 올립니다. 되감았으므로 그 뒤 read 도 같은 오류입니다 - zstd 모듈을 설치하거나 seek_to_offset 으로 그 배치를 넘기면 풀립니다(그 파티션을 seek 하면 미뤄 둔 오류도 지웁니다). 다른 파티션이 바빠도 오류가 묻히지 않습니다(레코드를 돌려준 다음 호출은 늘 오류입니다).
    • 분류: zstd 모듈 없음 / 모르는 코덱은 설치·설정 문제라 MIMQC_USER_FAULT / MIMQE_NOT_SUPPORTED, 깨진 자료는 예전처럼 MIMQC_FATAL_ERROR / MIMQE_INTERNAL_ERROR 입니다. 사유 문구(파티션, baseOffset, 되감은 자리 또는 되감기 거절 사유)는 get_report_msg() 와 원인(__cause__)에 있습니다.
    • 되감기 seek 가 거절되면(재배정 뒤 이 멤버 것이 아님 4161, 수동 배정 밖) 그대로 둡니다 - 새 주인이 커밋된 자리부터 다시 읽습니다. 오류는 같게 올리고 문구에 거절 사유를 적습니다.
    • COMMIT_IMMEDIATE 는 엔진이 읽는 순간 커밋했지만 되감으므로 이 멤버가 그 배치를 다시 받습니다(그 배치만 at-least-once). 다시 받기 전에 멤버가 끝나면 그 배치는 다시 오지 않을 수 있습니다(IMMEDIATE 는 at-most-once).
    • listen() 은 받은 레코드를 콜백한 뒤 on_error 로 알리고 멈춥니다(요청을 되풀이하며 돌지 않습니다). 패턴 구독은 그 토픽을 떼어 내지 않고 그 파티션만 되감으며, 같은 응답의 다른 토픽 레코드는 그대로 줍니다.
    • 패턴 구독 read_batch 는 요청 실패(요청 전체 거절 / 통신 실패)에서도 이 호출에서 이미 모은 레코드를 돌려주고 오류는 다음 호출이 올립니다(예전 판은 모은 레코드를 버렸는데 AUTO 커밋 좌표는 이미 병합돼 있어 앱이 못 본 레코드가 커밋됐습니다).
  • NIO 커넥터 송신은 기한 하나를 씁니다(TCPConnectorNIO, set_conn 으로 고른 경우만): 송신 한 번 전체가 접속 대기 시간 (최소 1 s)을 기한으로 씁니다. 기한이 지나면 MIMQE_SESSION_TIMEOUT(write waiting time has exceeded the time limit : [Nms] / written : [n / 전체])을 내고 연결을 닫습니다. 전에는 보낼 자리가 날 때까지 기다릴 때마다 한도를 새로 주어, 상대가 조금씩만 읽으면 송신이 끝없이 길어졌습니다. 기본(블로킹) 커넥터는 그대로입니다.
  • 교차 검증 도구 tests/cmpr_vectors.py: write <dir> 는 실제 producer 인코딩 경로로 만든 배치(py-none / py-gzip / py-zstd .batch)와 기대 레코드(.json)를 쓰고, check <dir> 는 디렉터리의 모든 *.batch 를 실제 소비 파싱 경로로 풀어 짝 .json 과 대조합니다. 자바 쪽 같은 도구가 쓴 java-*.batch 도 함께 검사합니다.

0.7.4 변경 - 패턴 토픽 떼어 내기 = 그룹 탈퇴(LEAVE), seek 규칙 (자바 v2.3.5 와 같음)

  • 엔진 요구: 아래 새 동작은 엔진 7.0.1.3417 이상(CONSGRP-1)에서 나옵니다. 그보다 옛 엔진에서는 그룹 탈퇴 요청을 조용히 건너뛰고(예전처럼 패턴 구독을 닫을 때까지 멤버로 남습니다), 파티션을 생략한 seek_to_offset() 은 0번 파티션을 되감습니다. 클라이언트는 엔진 판을 따로 검사하지 않습니다.
  • 패턴 구독 - 떼어 낸 토픽은 그룹에서 빠집니다: 재평가로 더 이상 매칭되지 않거나 read 중 오류로 떼어 낸 토픽은 그 (토픽, 구독) 그룹에서 이 세션을 뺍니다(MIMQ_TOPIC_LEAVE = 3401). 엔진이 곧바로 재배정하므로 같은 구독 이름의 다른 멤버가 그 파티션을 이어받고, 커서는 남습니다. 예전에는 연결을 닫지 않으니 패턴 구독을 닫을 때까지 멤버로 남아 그 파티션을 아무도 읽지 않았습니다. 그 토픽이 실린 다중 토픽 read 가 엔진에 걸려 있는 동안은 보내지 않고(엔진이 MIMQE_READ_IN_PROGRESS 로 거절합니다) 그 read 가 돌아온 직후 / 다음 read 전에 보냅니다 - 다른 스레드의 refresh_now() 는 걸린 read 를 기다리지 않고 돌아옵니다. best-effort 라 실패(구독 없음, 옛 엔진, 통신 실패)는 앱에 올리지 않습니다. 패턴 close() / unsubscribe() 와 단일 토픽 구독의 close() 는 그대로입니다(연결을 닫으면 엔진이 멤버를 뺍니다).
  • seek 는 이 멤버에게 배정된 파티션만: 배정 밖 파티션이면 ILOperationException (MIMQE_TOPIC_REBALANCED (partition N is not owned by this member)) 입니다. 재배정 통지가 아니라 거절이라 받아 둔 레코드는 그대로이고 다시 보내지 않습니다. 아직 read 하지 않은 구독은 그룹에 멤버가 하나도 없을 때만 seek 할 수 있습니다 - "subscribe → seek → read" 는 혼자일 때만 되고, 다른 멤버가 있으면 첫 read 로 배정 통지를 받은 뒤(예: set_rebalance_listener() 콜백 안) seek 하세요. assign() 구독은 배정 목록 밖 파티션이면 (Partition : N not in manual assignment) 사유로 거절됩니다. 관리 연결 (ILAdminTopicSubscription)의 seek 는 그룹에 멤버가 없을 때만 됩니다.
  • seek_to_offset(offset) - 파티션 생략은 엔진이 정합니다: 0.7.3 은 생략하면 ;partition:0 을 붙였습니다. 이제 offset:N 그대로 보냅니다. 새 엔진은 파티션이 하나인 토픽에서만 받고, 다중 파티션 토픽이면 MIMQE_INVALID_ARGUMENT (Partition : REQUIRED, Available : 0~N) 로 거절합니다 - 다중 파티션 토픽에서는 partition 을 주세요. 옛 엔진은 0번 파티션을 되감습니다. 성공하면 받아 둔 레코드는 0번 파티션 것만 버리고, 거절되면 아무것도 버리지 않습니다.
  • seek_to_time(epoch_ms, partition=-1) (seekToTime(epochMs, partition)): 파티션을 줄 수 있습니다(time:T;partition:P, 받아 둔 레코드도 P 것만 버립니다). 새 엔진은 이 멤버에게 배정된 파티션만 옮기고(파티션을 주면 그 파티션만), 돌려주는 값은 옮긴 파티션 중 번호가 가장 작은 파티션의 결과 오프셋입니다(옮긴 것이 없으면 -1). 옛 엔진은 파티션을 무시하고 구독 전체를 옮기며 0번 파티션의 결과를 줍니다 - 옛 엔진에서는 파티션을 주지 마세요.
  • 문서: ILAdminTopicSubscription.seek_to_offset() 의 "-1 이면 전 파티션" 은 틀린 설명이었습니다(엔진은 0번 파티션만 되감았습니다).
  • NIO 커넥터(TCPConnectorNIO)도 송수신이 실패하면 연결을 닫습니다(기본 커넥터와 같게): 전에는 큰 프레임을 받다가 시간이 다하면 나머지가 소켓에 남은 채 연결이 살아 있어 같은 연결의 다음 요청이 그 조각부터 읽었고(invalid transmission delimiter), 자동 확인 GET 이면 그 메시지를 잃었습니다. 상대가 프레임 중간에 끊으면 끝없이 돌던 것도 곧바로 MIMQE_NIO_ERROR 로 끝나고, 기한이 지났어도 소켓에 다 와 있는 프레임은 돌려주며, 보낼 자리가 없을 때(BlockingIOError) 곧바로 실패하던 송신은 기다렸다 이어 보냅니다. 기본 (블로킹) 커넥터에서도 프레임 중간에 시간이 다한 자동 확인 GET 의 메시지는 잃습니다(최대 한 번 전달) - 아주 큰 메시지는 GET 대기 시간 / 세션 시간 초과를 넉넉히 주거나 트랜잭션 GET(auto_commit=False 세션 + commit())을 쓰세요.

0.7.3 변경 - 구독마다 전용 연결, 소비 프로토콜 v5, 발행 / 소비 CPU 절감 (자바 v2.3.4 와 같음)

  • 엔진 요구: 토픽 구독은 엔진 7.0.1.3415 이상(소비 프로토콜 v5)에서만 동작합니다. 그보다 옛 엔진은 요청 태그를 모르므로 구독 연결을 끊습니다. 발행과 큐 기능은 영향이 없습니다.

  • 구독마다 전용 연결: 구독(패턴 구독이면 패턴 하나)마다 엔진 세션을 하나 따로 엽니다(클라이언트 이름 <이름>-SUB). 한 구독이 레코드를 기다려도 다른 구독이나 같은 ILQmgr 의 큐 / 관리 요청이 기다리지 않습니다. 대신 큐 관리자의 세션 수가 구독 수만큼 늘어납니다. 구독을 close() 하면 그 세션이 닫히고(엔진은 곧바로 멤버를 빼고 재배정합니다), ILQmgr.disconnect() 는 그 세션에서 만든 구독을 모두 닫으며, reconnect() 는 구독 세션도 함께 다시 붙입니다. close() 뒤의 read / commit 은 오류입니다. 다른 스레드의 read 가 레코드를 기다리던 중에 close() 하면 그 read 도 곧바로 끝납니다.

  • close() 와 unsubscribe() 의 차이: close() 는 이 멤버만 빠집니다(연결을 끊으면 엔진이 곧바로 멤버를 빼고 재배정하며, 커서는 남습니다). unsubscribe() 는 구독(그룹) 자체를 지웁니다 - 모든 파티션의 커서(durable 포함)가 사라지고, 같은 구독 이름으로 붙어 있는 다른 멤버(다른 프로세스 포함)도 다음 read 에서 MIMQE_TOPIC_SUBSCRIPTION_NOT_FOUND 를 받습니다. 구독을 없앨 때만 부르세요. unsubscribe() 는 close() 뒤에도 됩니다(ILQmgr 의 연결로 보냅니다). 구독 세션은 ILQmgr 의 세션 모드와 상관없이 늘 자동 확인 모드로 열립니다 - 토픽 read / commit / seek 는 원래 세션 트랜잭션에 묶이지 않으므로(qmgr.commit() / rollback() 은 큐만 다룹니다) 동작은 그대로입니다.

  • 소비 프로토콜 v5: 구독 연결의 요청에 태그를 붙여 겹쳐 보냅니다 - 레코드를 기다리는 동안에도 같은 구독의 commit / seek 가 곧바로 처리됩니다(한 구독의 read 는 하나씩 나갑니다). 패턴 구독은 매칭된 토픽 전체를 요청 한 번(다중 토픽 read)으로 받습니다. 커밋 좌표는 파티션 수와 관계없이 한 요청에 실리고, 한 요청 전체가 되거나 안 되거나입니다 - 이 멤버에게 배정되지 않은 파티션이 하나라도 있으면 MIMQE_TOPIC_REBALANCED (partition N is not owned by this member) 로 아무것도 커밋되지 않습니다.

  • 정확성:

    • 한 ILQmgr 에서 listen() 구독이 둘 이상이거나, listen() 중에 같은 ILQmgr 로 다른 요청(큐 put/get, 관리 요청, 다른 구독의 read)을 하면 소켓 하나에서 요청과 응답이 섞일 수 있었습니다. 이제 연결마다 잠금 하나로 요청-응답을 짝 단위로 줄 세웁니다(먼저 기다린 쪽이 먼저).
    • listen() 콜백이 예외를 던져 멈추면, 뒤이은 close() 가 처리하지 않은 레코드까지 AUTO 커밋하던 결함을 고쳤습니다. 이제 콜백이 정상으로 끝난 레코드만 커밋 대상이고, 콜백하지 못한 레코드는 같은 핸들의 다음 read / listen 이 다시 줍니다.
    • seek_to_offset(offset, partition) 이 다른 파티션의 받아 두기만 한 레코드와 AUTO 커밋 좌표까지 버리던 결함을 고쳤습니다(0.7.2 에도 있었습니다). 엔진은 지정한 파티션만 되감으므로 버려진 다른 파티션 구간은 다시 오지 않았고, AUTO 는 뒤 커밋이 그 구간을 건너뛰었습니다. 이제 지정한 파티션만 버리고, seek 응답보다 먼저 도착한(곧 seek 전 위치에서 읽은) 그 파티션의 레코드는 앱에 주지 않습니다. partition 을 생략하면 엔진은 0번 파티션 하나만 되감습니다(offset=0 이어도 구독 전체가 아닙니다). 예전에는 이때도 모든 파티션의 받아 둔 레코드를 버렸습니다 - 이제 ;partition:0 을 붙여 보내고 0번 것만 버립니다. seek_to_time() 은 전 파티션을 되감으므로 전부 버립니다.
    • seek 는 엔진이 멤버를 보지 않고 커서를 옮깁니다. 이 멤버에게 배정된 파티션만 되감고(get_assignment()), seek_to_time() 은 그룹에 멤버가 하나일 때만 쓰세요.
    • close() 의 AUTO 커밋이 MIMQE_TOPIC_REBALANCED (partition N is not owned by this member) 로 거절되면 그 파티션을 빼고 다시 보냅니다 - 아직 소유한 파티션의 좌표는 반영됩니다(예전에는 전부 버려져 다음 멤버에게 다시 갔습니다).
  • 문서: 세션 목록(ILSessionProperty)의 destroyedTime 은 살아 있는 세션에서 -1 입니다(0 이 아닙니다). 끊긴 세션은 끊긴 시각(epoch 밀리초)입니다. 살아 있는지는 status 로 판정하세요.

  • 발행 CPU: send() 가 레코드를 배치 버퍼에 곧바로 붙이고, 잠금 / 파티션 선택 / 키 인코딩 비용을 줄였습니다(레코드당 약 30% 감소).

  • 소비 CPU: 저장 v4 배치 응답을 풀면서 레코드 객체를 바로 만들고, AUTO 커밋 병합을 파티션마다 한 번 합니다(read_batch 약 40% 감소).

  • 소비 왕복:

    • read() 가 AUTO / MANUAL 에서 prefetch 건수(기본 500)까지 한 번에 받아 두고 차례로 줍니다(레코드마다 왕복하지 않습니다). IMMEDIATE 는 지금처럼 1건씩 요청합니다. 받아 두기만 한 레코드는 커밋되지 않으며, AUTO 커밋은 다음 요청(또는 close())에 실립니다.
    • commit(list) / close() 가 파티션마다 가장 큰 오프셋을 모아 파티션 수와 관계없이 한 요청으로 커밋합니다.
    • listen() 과 prefetch 기본 배치가 100 에서 500 이 됐습니다.
  • 패턴 구독: read / read_batch 한 번이 매칭된 토픽 전체를 다중 토픽 read 요청 하나로 묻습니다. 한 토픽이라도 레코드가 있으면 곧바로 돌아오고, 모두 비었으면 timeout_ms 동안 어느 토픽이든 들어오기를 기다립니다 - 한가한 토픽이 많아도 바쁜 토픽이 늦어지지 않습니다. 한 번에 토픽 하나가 가져가는 양은 bufferPerTopic 으로 묶고, 시작 토픽을 호출마다 한 칸씩 돌립니다(바쁜 토픽이 셋 이상일 때 하나를 건너뛰던 결함을 고쳤습니다). 목록 commit 은 토픽마다 한 요청이며, 보내기 전에 전부 검사합니다 (출처가 없는 레코드가 있으면 아무것도 커밋하지 않습니다).

0.7.2 변경 - batchSize 기본값 262144 로 되돌림, 발행 경로 CPU 절감 (자바 v2.3.2 와 같음)

  • batchSize 기본값을 262144 로 되돌렸습니다. 0.7.1 의 16384 에서 파이썬 클라는 처리량이 24~35% 낮았습니다(GIL 이 천장이라 배치마다 드는 비용이 그대로 처리량에서 빠집니다). 멱등 기본 켜짐은 그대로입니다.
  • 발행 경로 CPU 를 줄였습니다(API 는 그대로입니다). 레코드 future 가 배치 결과 하나를 함께 가리키고(완결은 배치당 한 번), 키 -> 파티션 결과를 캐시하고, send() 가 설정값을 레코드마다 다시 읽지 않고, I/O 스레드를 깨우는 신호를 합치고, ACK 하나의 배치 완결을 잠금 한 번에 묶습니다.
  • producer 는 만들 때 설정 객체를 복사합니다(자바 v2.3.2 도 같고, Kafka 와 같습니다). 만든 뒤 넘긴 객체를 바꿔도 그 producer 에는 반영되지 않습니다 - 설정을 바꾸려면 producer 를 새로 만드십시오.
  • 문서: 멱등이 거르는 것은 producer 안의 재전송뿐입니다 - 앱이 send 를 다시 부르면 새 레코드로 적재됩니다.

0.7.1 변경 - producer 기본값: 멱등 켜짐, batchSize 16384 (자바 v2.3.1 과 같음)

  • enableIdempotence 기본값이 켜짐입니다. 멱등을 직접 정하지 않았으면 acks(0) 이나 maxInFlight 가 5 를 넘을 때 멱등이 저절로 꺼집니다. enableIdempotence(True) 를 명시하고 그 둘과 함께 쓰면 전처럼 생성 때 ValueError 입니다.
  • 멱등이면 토픽을 비우거나(FLUSH) epoch 가 바뀐 순간 걸려 있던 배치를 다시 보내지 않고 실패로 올립니다 (MIMQE_TOPIC_FLUSHED / MIMQE_TOPIC_EPOCH_MISMATCH). 0.7.0 기본처럼 조용히 다시 보내려면 enableIdempotence(False) 를 명시하십시오.
  • batchSize 기본값이 16384 입니다(파티션 배치 하나의 상한). 0.7.0 은 262144 였습니다 - 처리량이 모자라면 batchSize(262144) 처럼 늘리십시오.
  • producer 수 안내: 한 토픽의 producer(연결)는 1~4개면 충분합니다(엔진 7.0.1.3407 실측).

0.7.0 변경 - 발행 v2: producer 하나 = 토픽 하나, 파티션별 배치, 봉투 겹쳐 보내기 (엔진 v2 필요)

발행 경로를 엔진의 새 발행 형식(v2)으로 바꿨습니다. 엔진도 v2 판이어야 합니다 - 옛 엔진에는 발행이 되지 않고 (MIMQE_NOT_SUPPORTED), 옛 클라이언트(0.6.x)는 v2 엔진에 발행할 수 없습니다. 둘을 함께 올리십시오.

  • producer 하나 = 토픽 하나 = 연결 하나. qmgr.create_producer("ORDER.EVENT", config) 로 만듭니다. 생성할 때 토픽의 파티션 수를 받아 두므로 없는 토픽이면 생성 단계에서 MIMQE_OBJECT_NOT_FOUND 로 실패합니다. 토픽이 여러 개면 토픽마다 producer 를 만듭니다. 한 토픽의 producer(연결)는 1~4개면 충분합니다 - 엔진 7.0.1.3407 실측에서 8개로 늘리면 어느 크기든 오히려 줄었습니다. 큰 레코드(1KB 급)는 1개로도 디스크가 먼저 한계에 닿고, 작은 레코드(100B 급)는 4개 근처가 최대입니다.
  • 레코드 토픽이 producer 토픽과 다르면 send() 가 MIMQE_INVALID_ARGUMENT 로 거부합니다.
  • 옛 create_producer(config) / ILTopicProducer(host, port, name, config) 모양으로 부르면 TypeError 가 납니다.
  • 배치는 파티션마다 모입니다. batchSize(기본 262144, 자바와 같음)는 파티션 배치 하나의 상한입니다. 보낼 때는 준비된 파티션 배치들을 봉투 하나에 maxRequestSize 까지 싣고, 봉투를 maxInFlight(기본 5) 개까지 겹쳐 보냅니다. 완결은 I/O 스레드가 하므로 lingerMs=0 이어도 send() 가 돌아온 시점에 끝나 있다는 보장은 없습니다 - 동기 발행은 send(rec).get() 입니다. send(rec, callback) 으로 콜백(callback(metadata, exc))을 걸 수 있습니다.
  • 멱등 발행은 파티션 배치 단위로 번호를 매깁니다. 멱등이면 maxInFlight 는 5 이하여야 합니다.
  • 버퍼(bufferMemory)가 차면 send() 가 maxBlockMs 만큼 기다렸다가 MIMQE_QUEUE_FULL (buffer memory full ...) 로 거부합니다(전에는 그 자리에서 비웠습니다). 부하 문제이므로 다시 시도할 수 있습니다.
  • 크기 상한: 레코드 하나가 혼자 실린 봉투가 maxRequestSize 이하, key 와 속성은 각각 99,999바이트까지입니다. 넘으면 send() 가 MIMQE_MSG_SIZE_OVER (...) 로 거부합니다. 배치는 토픽 세그먼트 크기에도 맞춰 끊습니다. 봉투 건수 상한(10만 건)은 없어졌습니다.
  • 멱등 발행에서 리더가 바뀌는 순간, 팔로워가 거절한 배치 뒤의 배치가 먼저 적재되면 앞 배치는 번호가 틈으로 남습니다. 이 배치는 MIMQE_OUT_OF_ORDER_SEQUENCE (sequence gap ...) 로 실패합니다(적재되지 않았습니다). 전에는 이 경우 성공으로 보고되고 레코드가 사라질 수 있었습니다.
  • 소비자 배치 읽기가 엔진 저장 v4 응답(2042)을 받습니다. 엔진은 배치를 쪼개지 않아 요청한 개수보다 많이 줄 수 있습니다. read_batch(max, t) 는 그래도 max 이하를 돌려주고, 넘친 레코드는 구독이 들고 있다가 다음 read / read_batch 가 먼저 줍니다. COMMIT_AUTO / COMMIT_MANUAL 은 앱에 건넨 레코드까지만 커밋하므로 들고 있던 레코드는 끊기면 다시 옵니다. COMMIT_IMMEDIATE 는 엔진이 읽는 순간 배치 끝까지 커밋하므로 들고 있던 레코드도 이미 커밋된 상태라, 앱에 건네기 전에 끊기면 다시 오지 않습니다(at-most-once). seek 는 들고 있던 레코드를 버립니다. 옛 배치 응답(2024)은 받지 않습니다 - 배치 읽기에는 저장 v4 엔진이 필요합니다.
  • 없어진 API: ILTopicFrameMsg 의 발행 조립(createPublish, buildBatchPublishEnvelope, patchPidSeq, setAckMode, setPartition, setPidSeq, PUBLISH_*_OFFSET)과 옛 응답 파싱(getCorrId, getResultCode, getReason, getResultPartition / Offset / Timestamp / Reason), 상수 ILC.MIMQ_PUBLISH_MESSAGE / MIMQ_PUBLISH_ACK_MESSAGE / MIMQ_BATCH_PUBLISH_MESSAGE / MIMQ_BATCH_PUBLISH_ACK_MESSAGE. ILTopicFrameMsg 는 구독 전달 수신 전용이 됐습니다.
  • 추가: ILTopicProducer.get_topic() / get_partition_count(), ILC.MIMQ_TOPIC_PUBLISH_V2_MESSAGE(2040) / MIMQ_TOPIC_PUBLISH_V2_ACK_MESSAGE(2041), 결과 코드 ILC.MIMQ_TOPIC_PUBLISH_DUPLICATE(4164) ~ FOLLOWER_NODE(4172), ILC.MIMQ_TOPIC_BATCH_V4_MESSAGE(2042)와 수신 전용 ILTopicBatchV4Msg.

0.6.11 변경 - 메시지 크기 상한(기본 128MiB), 재시도 보호, 클러스터 편의 개선

  • 메시지(프레임) 크기 상한의 기본값이 128MiB 입니다(엔진 Runtime@MaxFrameLength 기본과 같습니다). 넘는 메시지는 보내기 전에 MIMQE_MSG_SIZE_OVER 로 막고 연결은 그대로 둡니다.
    • 엔진은 서버의 상한을 알려 주지 않습니다. 서버에서 상한을 올렸다면 클라이언트도 같이 올리십시오 - 연결은 setMaxFrameLength(n)(ILQmgr / ILAdminService), producer 는 ILProducerConfig.maxFrameLength(n) 입니다. 범위는 1,048,576 ~ 999,999,999 입니다.
    • producer 레코드 한 건의 실제 상한은 maxRequestSize 와 maxFrameLength - 25(헤더) 중 작은 쪽입니다.
  • 재시도 보호: 1MiB 를 넘는 레코드를 보낸 직후 응답 없이 연결이 연속 2번 끊기면, producer 는 재시도를 멈추고 그 레코드를 MIMQE_MSG_SIZE_OVER 로 실패시킵니다. 서버 상한이 클라이언트 설정보다 작을 때 같은 큰 레코드를 시한이 다할 때까지 되풀이해 보내지 않게 하려는 것입니다. 큰 레코드는 따로 보내므로 옆 레코드는 영향을 받지 않습니다.
  • ILClusterNode("QM1", "10.0.0.12", 21001, 27091) 처럼 자바와 같은 순서로 만들 수 있습니다(ILClusterProperty("C1", 27091) 도 같습니다).
  • getAuthorityHolderType() 은 권한 보유자를 "MAIN" / "SUB" / "" 로 줍니다. getAuthorityHolder() 는 선 코드("M" / "S")라 clusterType 과 바로 비교하면 늘 거짓입니다.

★ 엔진 7.0.1.3404 전에는 엔진 쪽 상한이 999,999,999 였습니다. 그런 엔진에 128MiB 를 넘는 메시지를 보내려면 위 설정을 올리십시오.

0.6.10 변경 - 클러스터 관리 API 를 판사/조정자 모델로 (클러스터는 엔진 7.0.1.3391 이상)

클러스터 하나는 관리 서비스 셋(판사 · MAIN · SUB)에 각자 정의를 만듭니다. 한 서비스에 만들면 나머지에 퍼지던 옛 방식은 엔진에서 없어졌습니다.

  • 정의는 ILClusterProperty.judge(name, main, sub, ilcc_port) 와 ILClusterProperty.coordinator(name, "MAIN" 또는 "SUB", ilcc_port) 로 만들고, 세 서비스에서 각각 createCluster(정의) 합니다.
  • 노드(ILClusterNode)에 그 호스트 조정자의 ilcc_port 를 싣습니다. 판사 정의와 그 조정자 정의 두 곳에 같은 값을 적습니다.
  • setClusterProperty 로 바꿀 수 있는 것은 자동 기동(어디서나)과 하트비트(판사에서만)뿐입니다. 타입 · ilcc 포트 · 노드를 바꾸려면 세 곳에서 지우고 다시 만듭니다.
  • removeCluster / startCluster / stopCluster 는 접속한 서비스만 다룹니다. 떠 있어도 지워집니다. stopCluster 는 정지를 기다리지 않으므로, 곧바로 다시 띄우려면 상태가 STOPPED 인지 먼저 확인하십시오.
  • getClusterDefDiag 는 새 본문(FORMAT 2)을 읽습니다. 판사에 물으면 두 조정자의 상태가, 조정자에 물으면 판사와의 관계가 옵니다.
  • 기동/정지 실패 사유가 "-1" 대신 엔진이 준 사유로 올라옵니다.
  • 없어진 API: addClusterNode / removeClusterNode, getLastClusterWarning, 이름 · 포트로 만드는 createCluster 형태, ILClusterProperty 의 priSvc · secSvc · witSvc · 세 포트 · role · add_node · remove_node, 옛 진단 접근자(getLocal, getLocalAddress, isMissing, hasAgreement 등), 상수 CLUSTER_NO_NODES_MIN_REVISION · CLUSTER_DEFDIAG_MAX_FORMAT · CLUSTER_DEFDIAG_MAX_ELAPSED_MS.

그 밖:

  • 프레임 길이가 999,999,999 바이트를 넘으면 보내기 전에 MIMQE_MSG_SIZE_OVER 로 막습니다. 엔진은 그런 프레임을 받으면 사유 없이 연결을 끊습니다. maxRequestSize 는 999,999,974 를 넘지 않게 보정됩니다.

0.6.9 변경 - 파티션을 클라이언트가 정한다 (엔진 7.0.1.3401 이상 필요)

엔진이 더는 파티션을 고르지 않습니다. 이 판부터 producer 와 관리 API 단건 발행이 파티션을 정해 보냅니다(Kafka 와 같은 방식).

  • 키가 있으면 Kafka 와 같은 murmur2 로 UTF-8 키 바이트를 해시해 정합니다. 같은 키는 파이썬·자바 어느 클라이언트든 같은 파티션으로 갑니다.
  • 키가 없으면 producer 는 한 파티션에 batchSize 바이트씩 붙였다가 다음 파티션으로 넘깁니다(Kafka 3.3+ 균일 스티키). 관리 API 단건 발행은 토픽별 라운드로빈입니다.
  • 토픽의 파티션 수는 처음 보낼 때 한 번 묻습니다. 알 수 없으면(없는 토픽 등) send() 가 예외를 던집니다. acks=0 도 같습니다.
  • 멱등 발행(enableIdempotence(True))은 (토픽, 파티션)마다 번호를 매깁니다. 한 레코드가 거절돼도(너무 큼 등) 뒤 레코드는 계속 적재되고, 여러 건을 한 봉투에 묶어 보냅니다.
  • maxRequestSize(기본 1MB)는 레코드 한 건과 봉투 하나의 크기를 함께 묶습니다.

★ 그 이전 엔진에서는 발행이 모두 실패합니다(파티션 수를 물을 수 없다). 그런 엔진에는 0.6.8 을 쓰십시오.

0.2.0 신규 - pub/sub 토픽

큐에 더해 pub/sub 토픽을 지원합니다. 발행된 메시지는 retention 정책이 지울 때까지 보존되어 모든 구독자에게 각자의 커서로 전달됩니다(팬아웃). 큐와 달리 소비해도 사라지지 않습니다.

엔진 v7.1.1 rev 3265 이상이 필요합니다.

구독

from ilink.qmgr import ILQmgr
from ilink.topic import ILTopic, ILSubscribeOptions

qmgr = ILQmgr()
qmgr.connect("127.0.0.1", 19999, "order-svc", True)

topic = qmgr.access_topic("ORDER.EVENT")
sub = topic.subscribe("settlement", ILSubscribeOptions()
                      .start_mode(ILTopic.START_EARLIEST)
                      .commit_mode(ILTopic.COMMIT_MANUAL))

for msg in sub.read_batch(100, 3000):
    print(msg.get_offset(), msg.get_key(), msg.get_data_string())
    sub.commit(msg)

콜백(push)으로 받을 수도 있습니다.

sub.listen(lambda m: print(m.get_data_string()),
           lambda e: print("error:", e))
...
sub.stop_listening()

발행

발행은 producer 단일 경로입니다. producer 하나는 토픽 하나에 묶입니다(0.7.0). 동기 발행은 send().get()을 씁니다.

from ilink.producer import ILProducerConfig, ILProducerRecord

prod = qmgr.create_producer("ORDER.EVENT", ILProducerConfig())
meta = prod.send(ILProducerRecord("ORDER.EVENT", "k1", b"payload",
                                  properties={"trace-id": "abc"})).get()
print(meta.get_partition(), meta.get_offset())
prod.close()

와일드카드(패턴) 구독

패턴에 맞는 여러 토픽을 한꺼번에 구독합니다.

pat = qmgr.access_pattern("ORDER.*")
print(pat.resolve())                    # 지금 매칭되는 토픽 (구독 안 함)

ps = pat.subscribe("audit")
m = ps.read(3000)
print(m.get_topic_name(), m.get_data_string())

*는 구분자 .를 포함해 매칭합니다. ORDER.*는 ORDER.KR뿐 아니라 ORDER.KR.SUB도 잡습니다. 정규식이 아니라 글로브입니다. .은 리터럴이라 ORDER.*는 ORDERING을 잡지 않습니다.

토픽 관리

토픽 생성/삭제/속성변경은 관리 표면 전용입니다.

from ilink.admin import ILAdminService

svc = ILAdminService(); svc.connect("127.0.0.1", 9998)
adm = svc.accessAdminQmgr("QMGR1")
adm.createTopic("ORDER.EVENT", "partitions=3")
print(adm.getTopicList())

0.2.0 신규 - 클러스터(HA) 페일오버

클러스터 큐 관리자에 접속하면 후보 주소를 캐싱해 두었다가, 리더가 바뀌어도 따라갑니다.

qmgr.connect("10.0.0.1", 5000, "APP", True)   # 주소 하나면 됩니다
print(qmgr.get_cluster_addresses())           # ['10.0.0.1:5000', '10.0.0.2:5000']
qmgr.reconnect()                              # 새 리더를 찾아 재접속

# 첫 접속 시점의 장애까지 대비하려면 목록으로 (포트 인자 없음)
qmgr.connect("10.0.0.1:5000,10.0.0.2:5000", "APP", True)

접속에 성공하면 핸드셰이크 직후 나머지 노드 주소를 서버에서 받아 캐싱하므로 주소는 하나만 주면 됩니다. 팔로워를 지목해도 거부에 실린 리더 힌트를 따라 자동으로 리더에 붙으며, 이는 단일 주소든 목록이든 마찬가지입니다. 그때 알게 된 리더 주소는 후보 캐시에 들어가 목록으로 받은 주소와 똑같이 쓰이고, set_endpoint_cache_file() 로 경로를 지정해 뒀다면 파일에도 반영됩니다(이미 있으면 그대로 둡니다).

목록은 첫 접속 시점의 장애까지 대비할 때 씁니다 — 하나만 준 그 주소가 죽어 있으면 힌트를 줄 상대조차 없기 때문입니다.

0.4.4 변경 - 자동 재접속은 옵션이 아니라 기본 동작

캐싱된 endpoint 목록이 있으면(= 클러스터) 통신 장애로 실패한 put/get 을 페일오버 재접속 후 한 번 자동 재시도합니다. 켜고 끄는 설정은 없습니다 (set_auto_reconnect() 는 제거했습니다 — 아무 일도 하지 않는 세터를 남겨 두면 "껐는데 왜 재접속하냐"는 혼란만 생깁니다).

상황 동작
미커밋 트랜잭션 있음 예외를 올립니다. 앱이 reconnect() 후 트랜잭션을 처음부터 다시 수행
미커밋 트랜잭션 없음 (auto-commit 포함) 내부에서 재접속 후 1회 재시도
endpoint 목록 없음 (비클러스터/구엔진) 예외를 올립니다 - 갈 곳이 없습니다

재시도는 at-least-once 입니다. 서버가 처리한 뒤 응답이 유실된 시점에 재시도하면 중복 put 이 생길 수 있습니다. 재시도는 같은 msgId 로 나가므로 수신측에서 메시지 ID 로 걸러낼 수 있습니다. 미커밋 트랜잭션이 없을 때만 재시도하므로 트랜잭션 유실은 없습니다.

장애 전환(failover) 시 앱이 해야 할 일

리더가 죽으면 새 리더가 뽑힐 때까지 아무도 접속을 받지 않습니다. 실측(2노드 클러스터, 엔진 7.0.1.3341) 약 15~20초가 걸립니다. 그 사이의 재접속 시도는 정상적으로 실패하므로, 한 번 실패했다고 끝내지 말고 재시도해야 합니다.

세션 종류에 따라 앱이 할 일이 갈립니다.

세션 실패 시점에 미커밋 앱이 할 일
auto-commit 없음 아무것도 안 해도 됩니다. 그 연산을 다시 부르기만 하면 라이브러리가 재접속·재시도합니다
transacted 없음 위와 같습니다
transacted 있음 reconnect() 로 자리를 옮기고 트랜잭션을 처음부터 다시 수행해야 합니다

access_queue() 는 자동 재접속 대상이 아닙니다. 끊긴 세션에서 부르면 그대로 실패하니, 재접속에 성공한 뒤에 핸들을 다시 얻으십시오. 기존 핸들을 계속 쓰는 쪽은 스스로 복구됩니다.

(1) 트랜잭션 재수행이 필요 없는 경우 — auto-commit

import time
from ilink.qmgr import ILQmgr
from ilink.exception import ILException, ILSessionException, ILOperationException

q = ILQmgr()
q.set_endpoint_cache_file("/var/run/myapp/ilink-endpoints.txt")   # 재기동 대비(선택)
q.connect("10.0.0.1", 5000, "collector", True)                    # auto-commit
queue = q.access_queue("APP.EVENTS")                              # 핸들은 한 번만 얻는다

def publish(payload):
    """페일오버가 나도 이 함수는 그대로 둔다 - 라이브러리가 복구한다."""
    for attempt in range(1, 11):
        try:
            return queue.put(payload)          # 실패하면 내부에서 재접속 + 1회 재시도
        except ILOperationException:
            raise                              # 서버가 거절 - 재시도해도 같다
        except (ILException, ILSessionException):
            if attempt == 10:
                raise                          # 승격이 20초 넘게 안 끝났다
            time.sleep(3)                      # 리더 승격을 기다렸다가 다시

앱 코드에 reconnect() 가 없다는 점이 요지입니다. 같은 큐 핸들로 put 을 다시 부르기만 하면, 승격이 끝난 시점의 호출이 새 리더에 붙어 성공합니다(실측 20.6초에 복구).

(2) 트랜잭션 재수행이 필요한 경우 — transacted

import time
from ilink.qmgr import ILQmgr
from ilink.exception import ILException, ILSessionException, ILOperationException

MAX_RETRY = 10

def process_batch(records):
    """한 번의 트랜잭션 = 전부 커밋되거나 전부 없던 일이 된다."""
    q = ILQmgr()
    q.set_endpoint_cache_file("/var/run/myapp/ilink-endpoints.txt")
    q.connect("10.0.0.1", 5000, "billing-svc", False)      # transacted
    try:
        for attempt in range(1, MAX_RETRY + 1):
            try:
                queue = q.access_queue("APP.ORDERS")       # 재접속 뒤 핸들을 다시 얻는다
                for rec in records:
                    queue.put(rec)
                q.commit()                                 # 여기까지 와야 확정된다
                return
            except ILOperationException:
                q.rollback()                               # 서버 거절 - 재시도 무의미
                raise
            except (ILException, ILSessionException):
                if attempt == MAX_RETRY:
                    raise
                try:
                    q.reconnect()                          # 승격 전이면 여기서 또 실패한다
                except Exception:
                    time.sleep(3)                          # 기다렸다가 다음 회차에 다시
                # 루프 처음으로 -> 배치 전체를 다시 넣는다
    finally:
        q.disconnect()

읽어야 할 포인트 셋입니다.

  1. 재시도 단위는 트랜잭션 전체입니다. 실패 지점부터 이어붙이면 안 됩니다 — 앞서 넣은 것들도 커밋되지 않았으므로 함께 사라집니다(실측: 미커밋 1건 유실, 재수행 후 depth 2 정상).
  2. reconnect() 자체가 실패할 수 있습니다. 승격 전에 부르면 실패하는 게 정상이라, 루프 안에서 간격을 두고 다시 불러야 합니다(실측 2번째 시도, 16.5초에 성공).
  3. ILOperationException 과 통신 예외를 갈라 잡습니다. 전자는 서버가 판단해 거절한 것이라 재시도가 무의미하고, 후자만 재접속 대상입니다. 파이썬의 예외 클래스는 평면 구조라 ILSessionException 은 ILException 의 하위 타입이 아닙니다 — 통신 장애를 잡으려면 둘 다 적어야 합니다(자바는 모두 ILException 으로 감싸므로 하나면 됩니다).

앱이 재기동되는 경우

캐싱은 메모리에만 있으므로 프로세스가 죽으면 사라집니다. 설정에 남은 옛 리더 주소로 다시 붙을 때 이렇게 갈립니다.

옛 리더의 상태 결과
살아 있고 팔로워로 강등 그 노드가 리더 힌트를 주므로 자동으로 새 리더에 접속
죽어 있음 갈 곳이 없어 실패 — set_endpoint_cache_file() 을 쓰거나 주소를 목록으로 주십시오
q = ILQmgr()
q.set_endpoint_cache_file("/var/run/myapp/ilink-endpoints.txt")
q.connect("10.0.0.1", 5000, "APP", True)   # 이 주소가 죽어 있어도 파일 후보로 붙는다

0.4.0 신규 - 관리 연결 하나로 pub/sub

ILAdminService(관리 포트, 보통 9998) 연결만으로 토픽 발행·구독이 가능합니다. 큐 관리자 리스너 포트에 따로 붙지 않아도 되므로, 방화벽이 관리 포트만 열린 환경에서 쓸 수 있습니다.

엔진 v7.0.1 rev 3295 이상이 필요합니다.

from ilink.admin import ILAdminService
from ilink.exception import ILNoMsgException
from ilink.topic import ILSubscribeOptions, ILTopic

svc = ILAdminService()
svc.connect("127.0.0.1", 9998, "admin-app")
aq = svc.accessAdminQmgr("QM1")

# 발행 - (offset, timestamp, partition)
off, ts, part = aq.publish("ORDER.EVENT", b"payload", key="order-1")

# 구독
topic = aq.accessTopic("ORDER.EVENT")
sub = topic.subscribe("audit", ILSubscribeOptions()
                      .startMode(ILTopic.START_EARLIEST)
                      .commitMode(ILTopic.COMMIT_MANUAL))
try:
    while True:
        try:
            msg = sub.read(3000)
        except ILNoMsgException:
            break
        print(msg.get_offset(), msg.get_key(), msg.get_data_string())
        sub.commit(msg)
finally:
    sub.close()

패턴(와일드카드) 구독도 같은 연결로 됩니다.

ps = aq.accessPattern("ORDER.*").subscribe("audit")

제약: 관리 연결에는 배치가 없어 레코드 한 건에 한 번 왕복합니다. 멱등 발행도 지원되지 않습니다(acks=1 고정). 대량 처리나 중복 제거가 필요하면 리스너 포트에 ILQmgr 로 붙어 ILTopicProducer 를 쓰세요.

0.3.0 신규 - 허브-스포크 연결 전환 헬퍼

허브에 접속해 스포크를 찾고 연결을 전환하는 세 단계를 한 번에 처리하는 connectToSpoke()가 추가되었습니다. Java API에도 같은 이름으로 있습니다.

from ilink.admin import ILAdminService

for name in ("S48", "S85"):
    svc = ILAdminService.connectToSpoke("10.10.1.95", 9998, "ADMIN", name)
    try:
        print(name, [q.getName() for q in svc.getQmgrList()])
    finally:
        svc.disconnect()

# 키를 이미 알고 있으면 목록 조회를 건너뛴다 (키는 설정 파일에 저장되어 재기동해도 유지)
svc = ILAdminService.connectToSpokeByKey("10.10.1.95", 9998, "ADMIN", spoke_key)

연결 전환 후 그 연결은 해당 스포크에 직접 접속한 것으로 에뮬레이션됩니다. 따라서 전환은 연결당 한 번뿐이고, 다른 스포크로 가려면 새 연결이 필요합니다. connectToSpoke()가 그 반복을 담당합니다. 실패하면 스스로 연결을 닫으므로 소켓이 새지 않습니다.

TCP 연결 방향이 스포크 → 허브 한 방향뿐이라 스포크 쪽에 인바운드 포트를 열지 않고도 관리할 수 있습니다. 실측상 릴레이 오버헤드는 없었습니다(getQmgrList() 중앙값 릴레이 18.3ms vs 허브 직결 18.4ms).

0.2.1 신규 - Java API 옵션 표면 일치

발행/구독 옵션이 Java API와 같은 이름의 빌더로 정리되었습니다. 기존 snake 표기 (start_mode / commit_mode / linger_ms ...)도 그대로 쓸 수 있습니다.

from ilink.producer import ILProducerConfig

cfg = (ILProducerConfig()
       .acks(1).lingerMs(5).batchSize(32768)
       .maxRequestSize(1048576)      # 봉투 하나의 상한(레코드 1건도) - 넘는 레코드는 send()가 거부
       .bufferMemory(33554432)       # 미전송 누적 상한 - 차면 maxBlockMs 까지 기다린다
       .retries(3).retryBackoffMs(100)
       .deliveryTimeoutMs(120000)    # 재시도를 포함한 완결 시한
       .enableIdempotence(True))     # PID/시퀀스 중복 제거 (acks=1, maxInFlight 5 이하 필요)

prod = qmgr.create_producer("ORDER.EVENT", cfg)   # 0.7.0: producer 하나 = 토픽 하나
print(prod.get_producer_id(), prod.is_connected())

구독 옵션에 리밸런스 리스너/prefetch/수동 파티션이 추가되었고, 파티션을 직접 고르는 assign()이 생겼습니다.

sub = topic.subscribe("app1", ILSubscribeOptions()
                      .startMode(ILTopic.START_EARLIEST)
                      .commitMode(ILTopic.COMMIT_MANUAL)
                      .expiryMs(600000)          # 멤버 유휴 만료 10분
                      .prefetch(100)
                      .listener(on_rebalance))

sub = topic.assign("app1", [0, 2], ILTopic.START_EARLIEST)   # 수동 파티션 배정

주의 (0.2.0에서 올라올 때): 옵션 값은 이제 게터로 읽습니다. cfg.acks / opt.durable 은 빌더 메서드이므로 값이 필요하면 cfg.get_acks() / opt.is_durable() 을 쓰세요. cfg.acks = 0 같은 직접 대입은 그대로 동작합니다. 같은 이유로 ILClusterProperty.isAutoStart 와 ILSpokeProperty.isRunning 도 Java처럼 메서드가 되었습니다 (prop.isAutoStart()).

개발자 가이드 문서

공개 API 779개 전부에 한국어 docstring이 붙어 있습니다. 편집기에서 함수 위에 마우스를 올리거나 help()로 파라미터 타입·기본값·허용값·예외를 바로 확인할 수 있습니다.

help(qmgr.access_queue)
help(ILAdminQmgr.getStatSeries)

값 객체(ILQueueProperty 등)는 필드 목록이 클래스 docstring에 정리되어 있습니다. getX() / setX() 접근자는 필드 이름에서 자동으로 만들어지므로, 필드 목록이 곧 접근자 목록입니다.

help(ILQueueProperty)      # 필드 이름 / 타입 / 기본값 / 의미

Release files for jetstream-api 0.17.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for jetstream-api 0.17.0
File Size Uploaded
jetstream_api-0.17.0.tar.gz 399.0 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for jetstream-api 0.17.0
File Interpreter ABI Platform
jetstream_api-0.17.0-py3-none-any.whl Python 3 none any Details

Total release size: 757.1 kB

Release files / jetstream_api-0.17.0.tar.gz

Download URL jetstream_api-0.17.0.tar.gz
Size 399.0 kB
Tags Source
SHA-256 checksum
How to use checksums
e5b33b911b35c4de76006a70db9cf8ddd66717b8c542d31973e81f35563246b7
BLAKE2b-256 checksum
How to use checksums
12754f83e9f00aa99e9e9091985a7de692ff02e0bf29e1b8f7cbba30082839eb
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.13.4

Release files / jetstream_api-0.17.0-py3-none-any.whl

Download URL jetstream_api-0.17.0-py3-none-any.whl
Size 358.1 kB
Tags Python 3
SHA-256 checksum
How to use checksums
2c6a0d7a06de4b3fe4c9b6c08b9c8a392977c74399cb6a346c5eeadcf8a571e6
BLAKE2b-256 checksum
How to use checksums
c7ce5438f79878e16ceeb4167b746c176f1ec0588b7d89ffe1af20a3b014eb29
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.13.4

Release history Release notifications | RSS feed

This release

0.17.0 This release

2 release files

0.16.0

2 release files

0.15.0

2 release files

0.14.0

2 release files

0.13.0

2 release files

0.12.0

2 release files

0.11.0

2 release files

0.10.0

2 release files

0.9.0

2 release files

0.8.0

2 release files

0.7.4

2 release files

0.7.3

2 release files

0.7.2

2 release files

0.7.1

2 release files

0.7.0

2 release files

0.6.11

2 release files

0.6.10

2 release files

0.6.9

2 release files

0.6.8

2 release files

0.6.7

2 release files

0.6.6

2 release files

0.6.5

2 release files

0.6.4

2 release files

0.6.3

2 release files

0.6.2

2 release files

0.6.1

2 release files

0.6.0

2 release files

0.5.6

2 release files

0.5.5

2 release files

0.5.4

2 release files

0.5.3

2 release files

0.5.2

2 release files

0.5.1

2 release files

0.5.0

2 release files

0.4.5

2 release files

0.4.4

2 release files

0.4.3

2 release files

0.4.2

2 release files

0.4.1

2 release files

0.4.0

2 release files

0.3.0

2 release files

0.1.6

2 release files

0.1.5

2 release files

0.1.4

2 release files

0.1.3

2 release files

0.1.2

2 release files

0.1.1

2 release files

0.1.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page