Pulsar와 함께 filesystem offloader 사용하기

Pulsar와 함께 filesystem offloader 사용하기 (Use filesystem offloader with Pulsar)

이 장은 filesystem offloader를 설치·구성하고 Pulsar에서 사용하는 모든 단계를 안내해요. Hadoop 분산 파일 시스템(HDFS)이나 NFS 같은 파일시스템으로 데이터를 오프로딩해 장기·저렴한 저장소를 활용하고 싶을 때 필요해요.

출처: 문서

본문

설치 (Installation)

이 섹션은 filesystem offloader를 설치하는 방법을 설명해요.

사전 준비 (Prerequisite)

  • Pulsar: 2.4.2 이상 버전

단계 (Steps)

  1. Pulsar tarball 다운로드.
  2. Pulsar offloaders 패키지를 다운로드·압축 해제한 뒤 Pulsar 디렉터리에 offloaders로 복사해요. 계층형 저장소 offloader 설치 참고.

구성 (Configuration)

note BookKeeper에서 파일시스템으로 데이터를 오프로딩하기 전에 filesystem offloader 드라이버의 일부 속성을 구성해야 해요.

또한 filesystem offloader가 자동으로 실행되도록 구성하거나 수동으로 트리거할 수도 있어요.

filesystem offloader 드라이버 구성 (Configure filesystem offloader driver)

broker.conf 또는 standalone.conf 구성 파일에서 filesystem offloader 드라이버를 구성할 수 있어요.

  • HDFS
  • NFS

HDFS - 필수(Required) 구성

| 파라미터 | 설명 | 예시 값 | | managedLedgerOffloadDriver | Offloader 드라이버 이름(대소문자 구분 없음). | filesystem | | fileSystemURI | 연결 주소. 기본 Hadoop 분산 파일 시스템에 접근하는 URI예요. | hdfs://127.0.0.1:9000 | | offloadersDirectory | Offloader 디렉터리 | offloaders | | fileSystemProfilePath | Hadoop 프로파일 경로. 구성 파일은 Hadoop 프로파일 경로에 저장돼요. Hadoop 성능 튜닝을 위한 다양한 설정을 포함해요. | conf/filesystem_offload_core_site.xml |

HDFS - 선택(Optional) 구성

| 파라미터 | 설명 | 예시 값 | | managedLedgerMinLedgerRolloverTimeMinutes | 토픽의 원장 롤오버 사이 최소 시간.

참고: 프로덕션 환경에서는 이 파라미터를 설정하지 않는 것을 권장해요. | 10 | | managedLedgerMaxEntriesPerLedger | 롤오버 전에 원장에 추가할 최대 항목 수.

참고: 프로덕션 환경에서는 이 파라미터를 설정하지 않는 것을 권장해요. | 50000 |

NFS - 필수(Required) 구성

| 파라미터 | 설명 | 예시 값 | | managedLedgerOffloadDriver | Offloader 드라이버 이름(대소문자 구분 없음). | filesystem | | offloadersDirectory | Offloader 디렉터리 | offloaders | | fileSystemProfilePath | NFS 프로파일 경로. 구성 파일은 NFS 프로파일 경로에 저장돼요. 성능 튜닝을 위한 다양한 설정을 포함해요. | conf/filesystem_offload_core_site.xml |

NFS - 선택(Optional) 구성

| 파라미터 | 설명 | 예시 값 | | managedLedgerMinLedgerRolloverTimeMinutes | 토픽의 원장 롤오버 사이 최소 시간.

참고: 프로덕션 환경에서는 이 파라미터를 설정하지 않는 것을 권장해요. | 10 | | managedLedgerMaxEntriesPerLedger | 롤오버 전에 원장에 추가할 최대 항목 수.

참고: 프로덕션 환경에서는 이 파라미터를 설정하지 않는 것을 권장해요. | 50000 |

filesystem offloader를 자동으로 실행 (Run filesystem offloader automatically)

네임스페이스 정책을 구성하면 임계값에 도달했을 때 데이터를 자동으로 오프로딩하도록 할 수 있어요. 임계값은 토픽이 Pulsar 클러스터에 저장한 데이터 크기에 기반해요. 토픽 저장소가 임계값에 도달하면 오프로딩 작업이 자동으로 트리거돼요.

| 임계값 | 동작 | | > 0 | 토픽 저장소가 임계값에 도달하면 오프로딩 작업을 트리거해요. | | = 0 | 브로커가 가능한 한 빨리 데이터를 오프로딩하게 해요. | | < 0 | 자동 오프로딩 작업을 비활성화해요. |

자동 오프로딩은 새 세그먼트가 토픽 로그에 추가될 때 실행돼요. 네임스페이스에 임계값을 설정했지만 토픽에 생성되는 메시지가 적다면, 현재 세그먼트가 가득 차기 전까지 filesystem offloader는 동작하지 않아요.

pulsar-admin 같은 CLI 도구로 임계값을 구성할 수 있어요.

예시 (Example)

이 예시는 pulsar-admin으로 filesystem offloader 임계값을 10 MB로 설정해요.

pulsar-admin namespaces set-offload-threshold --size 10M my-tenant/my-namespace

tip pulsar-admin namespaces set-offload-threshold options 명령에 대한 자세한 내용(플래그, 설명, 기본값, 약어 포함)은 Pulsar admin docs를 참고해요.

filesystem offloader를 수동으로 실행 (Run filesystem offloader manually)

개별 토픽에 대해 다음 방법 중 하나로 filesystem offloader를 수동으로 트리거할 수 있어요.

  • REST 엔드포인트 사용.
  • CLI 도구(예: pulsar-admin) 사용.

CLI 도구로 filesystem offloader를 수동으로 트리거하려면 토픽에 대해 Pulsar 클러스터에 유지해야 하는 최대 데이터 양(임계값)을 지정해야 해요. Pulsar 클러스터의 토픽 데이터 크기가 이 임계값을 초과하면 임계값을 더 이상 초과하지 않을 때까지 토픽의 세그먼트가 파일시스템으로 오프로딩돼요. 더 오래된 세그먼트가 먼저 오프로딩돼요.

예시 (Example)

  • 이 예시는 pulsar-admin으로 filesystem offloader를 수동으로 실행해요.
pulsar-admin topics offload --size-threshold 10M persistent://my-tenant/my-namespace/topic1

출력 (Output)

Offload triggered for persistent://my-tenant/my-namespace/topic1 for messages before 2:0:-1

tip pulsar-admin topics offload options 명령에 대한 자세한 내용(플래그, 설명, 기본값, 약어 포함)은 Pulsar admin docs를 참고해요.

  • 이 예시는 pulsar-admin으로 filesystem offloader 상태를 확인해요.
pulsar-admin topics offload-status persistent://my-tenant/my-namespace/topic1

출력 (Output)

Offload is currently running

파일시스템이 작업을 완료할 때까지 기다리려면 -w 플래그를 추가해요.

pulsar-admin topics offload-status -w persistent://my-tenant/my-namespace/topic1

출력 (Output)

Offload was a success

오프로딩 작업에 오류가 있으면 오류는 pulsar-admin topics offload-status 명령으로 전파돼요.

pulsar-admin topics offload-status persistent://my-tenant/my-namespace/topic1

출력 (Output)

Error in offload
null
Reason: Error offloading: org.apache.bookkeeper.mledger.ManagedLedgerException: java.util.concurrent.CompletionException: com.amazonaws.services.s3.model.AmazonS3Exception: Anonymous users cannot initiate multipart uploads.  Please authenticate. (Service: Amazon S3; Status Code: 403; Error Code: AccessDenied; Request ID: 798758DE3F1776DF; S3 Extended Request ID: dhBFz/lZm1oiG/oBEepeNlhrtsDlzoOhocuYMpKihQGXe6EG8puRGOkK6UwqzVrMXTWBxxHcS+g=), S3 Extended Request ID: dhBFz/lZm1oiG/oBEepeNlhrtsDlzoOhocuYMpKihQGXe6EG8puRGOkK6UwqzVrMXTWBxxHcS+g=

tip pulsar-admin topics offload-status options 명령에 대한 자세한 내용(플래그, 설명, 기본값, 약어 포함)은 Pulsar admin docs를 참고해요.

튜토리얼 (Tutorial)

이 섹션은 filesystem offloader를 사용해 Pulsar에서 HDFS(Hadoop Distributed File System) 또는 NFS(Network File System)로 데이터를 이동하는 단계별 지침을 제공해요.

HDFS로 데이터 오프로딩 (Offload data to HDFS)

tip 이 튜토리얼은 Hadoop 단일 노드 클러스터를 설정하고 Hadoop 3.2.1을 사용해요. Hadoop 단일 노드 클러스터 설정 방법은 여기를 참고해요.

1단계: HDFS 환경 준비 (Step 1: Prepare the HDFS environment)
  1. Hadoop 3.2.1을 다운로드·압축 해제해요.
wget https://mirrors.bfsu.edu.cn/apache/hadoop/common/hadoop-3.2.1/hadoop-3.2.1.tar.gz
tar -zxvf hadoop-3.2.1.tar.gz -C $HADOOP_HOME
  1. Hadoop을 구성해요.
# $HADOOP_HOME/etc/hadoop/core-site.xml
<configuration>
    <property>
        <name>fs.defaultFS</name>
        <value>hdfs://localhost:9000</value>
    </property>
</configuration>
# $HADOOP_HOME/etc/hadoop/hdfs-site.xml
<configuration>
    <property>
        <name>dfs.replication</name>
        <value>1</value>
    </property>
</configuration>
  1. 패스프레이즈 없는(passphraseless) ssh를 설정해요.
# Now check that you can ssh to the localhost without a passphrase:
ssh localhost
# If you cannot ssh to localhost without a passphrase, execute the following commands
ssh-keygen -t rsa -P '' -f ~/.ssh/id_rsa
cat ~/.ssh/id_rsa.pub >> ~/.ssh/authorized_keys
chmod 0600 ~/.ssh/authorized_keys
  1. HDFS를 시작해요.
# don't execute this command repeatedly, repeat execute will cause the clusterId of the datanode is not consistent with namenode
$HADOOP_HOME/bin/hadoop namenode -format
$HADOOP_HOME/sbin/start-dfs.sh
  1. HDFS 웹사이트로 이동해요. Overview 페이지를 볼 수 있어요.

  2. 상단 내비게이션 바에서 Datanodes를 클릭해 DataNode 정보를 확인해요.

  3. HTTP Address를 클릭해 localhost:9866에 대한 더 자세한 정보를 얻어요. 아래에서 볼 수 있듯이 Capacity Used의 크기는 4 KB로, 초기 값이에요.

2단계: filesystem offloader 설치 (Step 2: Install the filesystem offloader)

자세한 내용은 설치를 참고해요.

3단계: filesystem offloader 구성 (Step 3: Configure the filesystem offloader)

구성 섹션에서 언급했듯이, 사용하기 전에 filesystem offloader 드라이버의 일부 속성을 구성해야 해요. 이 튜토리얼은 아래와 같이 filesystem offloader 드라이버를 구성하고 Pulsar를 standalone 모드로 실행한다고 가정해요.

conf/standalone.conf 파일에 다음 구성을 설정해요.

managedLedgerOffloadDriver=filesystem
fileSystemURI=hdfs://127.0.0.1:9000
fileSystemProfilePath=conf/filesystem_offload_core_site.xml

note 테스트 목적으로 다음 두 구성을 설정해 원장 롤오버를 빠르게 할 수 있지만, 프로덕션 환경에서는 설정하지 않는 것을 권장해요.

managedLedgerMinLedgerRolloverTimeMinutes=1
managedLedgerMaxEntriesPerLedger=100
4단계: BookKeeper에서 파일시스템으로 데이터 오프로딩 (Step 4: Offload data from BookKeeper to filesystem)

Pulsar tarball을 다운로드한 저장소에서 다음 명령을 실행해요. 예: ~/path/to/apache-pulsar-2.5.1.

  1. Pulsar standalone을 시작해요.
bin/pulsar standalone -a 127.0.0.1
  1. 생성된 데이터가 즉시 삭제되지 않도록 보존 정책(retention policy)을 설정하는 것을 권장해요. 보존 정책은 크기(size) 제한 또는 시간(time) 제한이 될 수 있어요. 보존 정책에 더 큰 값을 설정할수록 데이터를 더 오래 유지할 수 있어요.
bin/pulsar-admin namespaces set-retention public/default --size 100M --time 2d

tip pulsarctl namespaces set-retention options 명령에 대한 자세한 내용(플래그, 설명, 기본값, 약어 포함)은 여기를 참고해요.

  1. pulsar-client로 데이터를 생성해요.
bin/pulsar-client produce -m "Hello FileSystem Offloader" -n 1000 public/default/fs-test
  1. 원장 롤오버가 트리거된 후 오프로딩 작업이 시작돼요. 오프로딩이 성공하려면 여러 원장 롤오버가 트리거될 때까지 기다리는 것을 권장해요. 이 경우 몇 초 기다려야 할 수 있어요. pulsarctl로 원장 상태를 확인할 수 있어요.
bin/pulsar-admin topics stats-internal public/default/fs-test

출력 (Output)

원장 696의 데이터는 오프로딩되지 않았어요.

{
"version": 1,
"creationDate": "2020-06-16T21:46:25.807+08:00",
"modificationDate": "2020-06-16T21:46:25.821+08:00",
"ledgers": [
{
    "ledgerId": 696,
    "isOffloaded": false
}],
"cursors": {}
}
  1. 몇 초 기다린 뒤 토픽에 더 많은 메시지를 보내요.
bin/pulsar-client produce -m "Hello FileSystem Offloader" -n 1000 public/default/fs-test
  1. pulsarctl로 원장 상태를 확인해요.
bin/pulsar-admin topics stats-internal public/default/fs-test

출력 (Output)

원장 696이 롤오버됐어요.

{
"version": 2,
"creationDate": "2020-06-16T21:46:25.807+08:00",
"modificationDate": "2020-06-16T21:48:52.288+08:00",
"ledgers": [
{
    "ledgerId": 696,
    "entries": 1001,
    "size": 81695,
    "isOffloaded": false
},
{
    "ledgerId": 697,
    "isOffloaded": false
}],
"cursors": {}
}
  1. pulsarctl로 오프로딩 작업을 수동으로 트리거해요.
bin/pulsar-admin topics offload -s 0 public/default/fs-test

출력 (Output)

원장 697 이전의 데이터가 오프로딩됐어요.

# offload info, the ledgers before 697 will be offloaded
Offload triggered for persistent://public/default/fs-test3 for messages before 697:0:-1
  1. pulsarctl로 원장 상태를 확인해요.
bin/pulsar-admin topics stats-internal public/default/fs-test

출력 (Output)

원장 696의 데이터가 오프로딩됐어요.

{
"version": 4,
"creationDate": "2020-06-16T21:46:25.807+08:00",
"modificationDate": "2020-06-16T21:52:13.25+08:00",
"ledgers": [
{
    "ledgerId": 696,
    "entries": 1001,
    "size": 81695,
    "isOffloaded": true
},
{
    "ledgerId": 697,
    "isOffloaded": false
}],
"cursors": {}
}

그리고 Capacity Used는 4 KB에서 116.46 KB로 바뀌었어요.

NFS로 데이터 오프로딩 (Offload data to NFS)

note 이 섹션에서는 NFS 서비스를 활성화하고 NFS 서비스의 공유 경로를 설정했다고 가정해요. 이 섹션에서는 /Users/test를 NFS 서비스의 공유 경로로 사용해요.

1단계: filesystem offloader 설치 (Step 1: Install the filesystem offloader)

자세한 내용은 설치를 참고해요.

2단계: NFS를 로컬 파일시스템에 마운트 (Step 2: Mount your NFS to your local filesystem)

이 예시는 /Users/pulsar_nfs/Users/test에 마운트해요.

mount -e 192.168.0.103:/Users/test /Users/pulsar_nfs
3단계: filesystem offloader 드라이버 구성 (Step 3: Configure the filesystem offloader driver)

구성 섹션에서 언급했듯이, 사용하기 전에 filesystem offloader 드라이버의 일부 속성을 구성해야 해요. 이 튜토리얼은 아래와 같이 filesystem offloader 드라이버를 구성하고 Pulsar를 standalone 모드로 실행한다고 가정해요.

  1. conf/standalone.conf 파일에 다음 구성을 설정해요.
managedLedgerOffloadDriver=filesystem
fileSystemProfilePath=conf/filesystem_offload_core_site.xml
  1. filesystem_offload_core_site.xml을 다음과 같이 수정해요.
<property>
    <name>fs.defaultFS</name>
    <value>file:///</value>
</property>
<property>
    <name>hadoop.tmp.dir</name>
    <value>file:///Users/pulsar_nfs</value>
</property>
<property>
    <name>io.file.buffer.size</name>
    <value>4096</value>
</property>
<property>
    <name>io.seqfile.compress.blocksize</name>
    <value>1000000</value>
</property>
<property>
    <name>io.seqfile.compression.type</name>
    <value>BLOCK</value>
</property>
<property>
    <name>io.map.index.interval</name>
    <value>128</value>
</property>
4단계: BookKeeper에서 파일시스템으로 데이터 오프로딩 (Step 4: Offload data from BookKeeper to filesystem)

HDFS로 데이터 오프로딩의 4단계를 참고해요.

파일시스템에서 오프로딩된 데이터 읽기 (Read offloaded data from filesystem)

  • 오프로딩된 데이터는 파일시스템의 다음 새 경로에 MapFile로 저장돼요.
path = storageBasePath + "/" + managedLedgerName + "/" + ledgerId + "-" + uuid.toString();
  • storageBasePathbroker.conf 또는 filesystem_offload_core_site.xml에 구성된 hadoop.tmp.dir의 값이에요.
  • managedLedgerName은 persistentTopic 관리자의 원장 이름이에요.
managedLedgerName of persistent://public/default/topics-name is public/default/persistent/topics-name.

다음 방법으로 managedLedgerName을 얻을 수 있어요.

String managedLedgerName = TopicName.get("persistent://public/default/topics-name").getPersistenceNamingEncoding();

파일시스템에서 데이터를 원장 항목으로 읽으려면 다음 단계를 완료해요.

  1. 새 경로의 MapFile과 파일시스템의 configuration 모두를 읽는 리더를 만들어요.
MapFile.Reader reader = new MapFile.Reader(new Path(dataFilePath),  configuration);
  1. 파일시스템에서 LedgerEntry로 데이터를 읽어요.
  LongWritable key = new LongWritable();
  BytesWritable value = new BytesWritable();
  key.set(nextExpectedId - 1);
  reader.seek(key);
  reader.next(key, value);
  int length = value.getLength();
  long entryId = key.get();
  ByteBuf buf = PooledByteBufAllocator.DEFAULT.buffer(length, length);
  buf.writeBytes(value.copyBytes());
  LedgerEntryImpl ledgerEntry = LedgerEntryImpl.create(ledgerId, entryId, length, buf);
  1. LedgerEntryMessage로 역직렬화해요.
     ByteBuf metadataAndPayload = ledgerEntry.getDataBuffer();
     long totalSize = metadataAndPayload.readableBytes();
     BrokerEntryMetadata brokerEntryMetadata = Commands.peekBrokerEntryMetadataIfExist(metadataAndPayload);
     MessageMetadata metadata = Commands.parseMessageMetadata(metadataAndPayload);
     Map<String, String> properties = new TreeMap();
     properties.put("X-Pulsar-batch-size", String.valueOf(totalSize
             - metadata.getSerializedSize()));
     properties.put("TOTAL-CHUNKS", Integer.toString(metadata.getNumChunksFromMsg()));
     properties.put("CHUNK-ID", Integer.toString(metadata.getChunkId()));
     // Decode if needed
     CompressionCodec codec = CompressionCodecProvider.getCompressionCodec(metadata.getCompression());
     ByteBuf uncompressedPayload = codec.decode(metadataAndPayload, metadata.getUncompressedSize());
     // Copy into a heap buffer for output stream compatibility
     ByteBuf data = PulsarByteBufAllocator.DEFAULT.heapBuffer(uncompressedPayload.readableBytes(),
             uncompressedPayload.readableBytes());
     data.writeBytes(uncompressedPayload);
     uncompressedPayload.release();
     MessageImpl message = new MessageImpl(topic, ((PositionImpl)ledgerEntry.getPosition()).toString(), properties,
             data, Schema.BYTES, metadata);
     message.setBrokerEntryMetadata(brokerEntryMetadata);

더 알아보기 (Learn more)