Azure Table Storage

Azure Table Storage

이 문서는 HadoopInputFormat 래퍼를 사용하여 기존 Hadoop 입력 포맷 구현으로 Azure의 Table Storage에 접근하는 예시입니다. 별도의 azure-tables-hadoop 프로젝트를 빌드하고 Flink 프로젝트에 의존성을 추가한 뒤, Hadoop 입력 포맷 래퍼로 Azure 테이블 데이터를 DataStream으로 읽어옵니다.

출처: 문서

본문

이 예시는 Azure의 Table Storage에 접근하기 위해 기존 Hadoop 입력 포맷 구현을 사용하는 HadoopInputFormat 래퍼를 사용합니다.

  1. azure-tables-hadoop 프로젝트를 다운로드하고 컴파일합니다. 이 프로젝트가 개발한 입력 포맷은 아직 Maven Central에 없으므로 프로젝트를 직접 빌드해야 합니다. 다음 명령을 실행하세요:
git clone https://github.com/mooso/azure-tables-hadoop.git
cd azure-tables-hadoop
mvn clean install
  1. quickstart를 사용하여 새 Flink 프로젝트를 설정합니다:
curl https://flink.apache.org/q/quickstart.sh | bash
  1. pom.xml 파일의 <dependencies> 섹션에 다음 의존성을 추가합니다:
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-hadoop-compatibility</artifactId>
    <version>{{< version >}}</version>
</dependency>
<dependency>
    <groupId>com.microsoft.hadoop</groupId>
    <artifactId>microsoft-hadoop-azure</artifactId>
    <version>0.0.5</version>
</dependency>

flink-hadoop-compatibility는 Hadoop 입력 포맷 래퍼를 제공하는 Flink 패키지입니다. microsoft-hadoop-azure는 이전에 우리가 빌드한 프로젝트를 우리 프로젝트에 추가합니다.

이제 프로젝트는 코딩을 시작할 준비가 되었습니다. 프로젝트를 IntelliJ 같은 IDE에 가져오는 것을 권장합니다. Maven 프로젝트로 가져와야 합니다. Job.java 파일을 찾아보세요. 이는 Flink 작업의 빈 스켈레톤입니다.

다음 코드를 붙여넣으세요:

import java.util.Map;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.java.DataStream;
import org.apache.flink.api.java.ExecutionEnvironment;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.hadoopcompatibility.mapreduce.HadoopInputFormat;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import com.microsoft.hadoop.azure.AzureTableConfiguration;
import com.microsoft.hadoop.azure.AzureTableInputFormat;
import com.microsoft.hadoop.azure.WritableEntity;
import com.microsoft.windowsazure.storage.table.EntityProperty;

public class AzureTableExample {

  public static void main(String[] args) throws Exception {
    // set up the execution environment
    final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

    env.setRuntimeMode(RuntimeExecutionMode.BATCH);
    // create a  AzureTableInputFormat, using a Hadoop input format wrapper
    HadoopInputFormat<Text, WritableEntity> hdIf = new HadoopInputFormat<Text, WritableEntity>(new AzureTableInputFormat(), Text.class, WritableEntity.class, new Job());

    // set the Account URI, something like: https://apacheflink.table.core.windows.net
    hdIf.getConfiguration().set(azuretableconfiguration.Keys.ACCOUNT_URI.getKey(), "TODO");
    // set the secret storage key here
    hdIf.getConfiguration().set(AzureTableConfiguration.Keys.STORAGE_KEY.getKey(), "TODO");
    // set the table name here
    hdIf.getConfiguration().set(AzureTableConfiguration.Keys.TABLE_NAME.getKey(), "TODO");

    DataStream<Tuple2<Text, WritableEntity>> input = env.createInput(hdIf);
    // a little example how to use the data in a mapper.
    DataStream<String> fin = input.map(new MapFunction<Tuple2<Text,WritableEntity>, String>() {
      @Override
      public String map(Tuple2<Text, WritableEntity> arg0) throws Exception {
        System.err.println("--------------------------------\nKey = "+arg0.f0);
        WritableEntity we = arg0.f1;

        for(Map.Entry<String, EntityProperty> prop : we.getProperties().entrySet()) {
          System.err.println("key="+prop.getKey() + " ; value (asString)="+prop.getValue().getValueAsString());
        }

        return arg0.f0.toString();
      }
    });

    // emit result (this works only locally)
    fin.print();

    // execute program
    env.execute("Azure Example");
  }
}

이 예시는 Azure 테이블에 접근하여 데이터를 Flink의 DataStream(보다 구체적으로 집합의 타입은 DataStream<Tuple2<Text, WritableEntity>>)으로 바꾸는 방법을 보여줍니다. DataStream을 사용하면 DataStream에 알려진 모든 변환을 적용할 수 있습니다.

더 알아보기 (Learn more)