Apache Cassandra 커넥터

Apache Cassandra 커넥터

이 커넥터는 Apache Cassandra 데이터베이스에 데이터를 쓰는 sink를 제공합니다.

이 커넥터를 사용하려면 프로젝트에 다음 의존성을 추가하세요.

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-cassandra_2.12</artifactId>
    <version>2.3.0</version>
</dependency>

스트리밍 커넥터는 현재 바이너리 배포판에 포함되지 않습니다. 클러스터 실행을 위해 링크하는 방법은 여기를 참조하세요.

출처: 문서

본문

Apache Cassandra 설치

로컬 머신에서 Cassandra 인스턴스를 구동하는 방법은 여러 가지가 있습니다.

  1. Cassandra Getting Started page의 지침을 따릅니다.
  2. 공식 Docker Repository에서 Cassandra를 실행하는 컨테이너를 구동합니다.

Cassandra 소스

Flink는 FLIP-27 유계(bounded) 소스를 제공하여 Cassandra에서 읽고 엔티티 컬렉션을 DataStream<Entity>로 반환합니다. 엔티티는 주석을 포함하는 POJO( Cassandra object mapper에 설명됨)를 기반으로 Cassandra 매퍼(MappingManager)에 의해 생성됩니다.

소스를 사용하려면 다음을 수행하세요.

ClusterBuilder clusterBuilder = new ClusterBuilder() {    
    @Override    
    protected Cluster buildCluster(Cluster.Builder builder) {      
        return builder.addContactPointsWithPorts(new InetSocketAddress(HOST,PORT))
                      .withQueryOptions(new QueryOptions().setConsistencyLevel(CL))                    
                      .withSocketOptions(new SocketOptions()                   
                      .setConnectTimeoutMillis(CONNECT_TIMEOUT)                    
                      .setReadTimeoutMillis(READ_TIMEOUT))                    
                      .build();    
    }
};  
long maxSplitMemorySize = ... //optional max split size in bytes minimum is 10MB. If not set, maxSplitMemorySize = 64 MB
Source cassandraSource = new CassandraSource(clusterBuilder, 
                                             maxSplitMemorySize, 
                                             Pojo.class, 
                                             "select ... from KEYSPACE.TABLE ...;",
                                             () -> new Mapper.Option[] {Mapper.Option.saveNullFields(true)});    
DataStream<Pojo> stream = env.fromSource(cassandraSource, WatermarkStrategy.noWatermarks(),  "CassandraSource");

성능과 관련하여 소스는 다음과 같이 테이블 데이터를 분할합니다: numSplits = tableSize/maxSplitMemorySize.

tableSize를 결정할 수 없거나 이전 numSplits 계산이 너무 적은 분할을 만들면 numSplits = parallelism으로 폴백합니다.

Cassandra Sink

구성

Flink의 Cassandra sink는 정적 CassandraSink.addSink(DataStream input) 메서드로 생성합니다. 이 메서드는 sink를 추가로 구성하는 메서드를 제공하는 CassandraSinkBuilder를 반환하고, 마지막으로 build()로 sink 인스턴스를 구성합니다.

다음 구성 메서드를 사용할 수 있습니다.

  1. setQuery(String query)
    • sink가 받는 모든 레코드에 대해 실행되는 upsert 쿼리를 설정합니다.
    • 쿼리는 내부적으로 CQL 문으로 처리됩니다.
    • Tuple 데이터 타입을 처리할 때는 반드시 upsert 쿼리를 설정하세요.
    • POJO 데이터 타입을 처리할 때는 쿼리를 설정하지 마세요.
  2. setClusterBuilder(ClusterBuilder clusterBuilder)
    • 일관성 수준, 재시도 정책 등 더 정교한 설정으로 cassandra 연결을 구성하는 데 사용되는 클러스터 빌더를 설정합니다.
  3. setHost(String host[, int port])
    • Cassandra 인스턴스에 연결할 host/port 정보가 있는 setClusterBuilder()의 간단한 버전입니다.
  4. setMapperOptions(MapperOptions options)
    • DataStax ObjectMapper를 구성하는 데 사용되는 매퍼 옵션을 설정합니다.
    • POJO 데이터 타입을 처리할 때만 적용됩니다.
  5. setMaxConcurrentRequests(int maxConcurrentRequests, Duration timeout)
    • 실행 허가를 획득하기 위한 타임아웃과 함께 허용되는 최대 동시 요청 수를 설정합니다.
    • **enableWriteAheadLog()**가 구성되지 않은 경우에만 적용됩니다.
  6. enableWriteAheadLog([CheckpointCommitter committer])
    • 선택 설정입니다.
    • 비결정적 알고리즘에 대한 exactly-once 처리를 허용합니다.
  7. setFailureHandler([CassandraFailureHandler failureHandler])
    • 선택 설정입니다.
    • 사용자 지정 실패 핸들러를 설정합니다.
  8. setDefaultKeyspace(String keyspace)
    • 사용할 기본 keyspace를 설정합니다.
  9. enableIgnoreNullFields()
    • null 값 무시를 활성화합니다. null 값을 unset으로 취급하고 null 필드 쓰기와 tombstone 생성을 피합니다.
  10. build()
    • 구성을 마무리하고 CassandraSink 인스턴스를 구성합니다.

Write-ahead Log

checkpoint 커미터는 완료된 checkpoint에 대한 추가 정보를 일부 리소스에 저장합니다. 이 정보는 실패 시 마지막 완료 checkpoint의 전체 재실행을 방지하는 데 사용됩니다. CassandraCommitter를 사용하여 이 정보를 cassandra의 별도 테이블에 저장할 수 있습니다. 이 테이블은 Flink가 정리하지 않을 것임에 유의하세요.

쿼리가 멱등(idempotent, 결과를 바꾸지 않고 여러 번 적용 가능)이고 checkpointing이 활성화되면 Flink는 exactly-once 보장을 제공할 수 있습니다. 실패 시 실패한 checkpoint가 완전히 재실행됩니다.

또한 비결정적 프로그램의 경우 write-ahead log를 활성화해야 합니다. 이러한 프로그램의 경우 재실행된 checkpoint는 이전 시도와 완전히 다를 수 있으며, 첫 번째 시도의 일부가 이미 쓰여졌을 수 있으므로 데이터베이스가 일관되지 않은 상태로 남을 수 있습니다. write-ahead log는 재실행된 checkpoint가 첫 번째 시도와 동일함을 보장합니다. 이 기능을 활성화하면 지연 시간에 부정적인 영향이 있을 수 있습니다.

참고: write-ahead log 기능은 현재 실험적입니다. 많은 경우 이를 활성화하지 않고 커넥터를 사용하는 것으로 충분합니다. 문제가 있으면 개발 메일링 리스트에 보고하세요.

Checkpointing 및 장애 허용

checkpointing이 활성화되면 Cassandra Sink는 C* 인스턴스에 대한 동작 요청의 at-least-once 전달을 보장합니다.

자세한 내용은 checkpoints 문서fault tolerance guarantee 문서를 참조하세요.

예시

Cassandra sink는 현재 Tuple과 POJO 데이터 타입을 모두 지원하며, Flink가 사용되는 입력 타입을 자동으로 감지합니다. 이러한 스트리밍 데이터 타입의 일반적인 사용은 Supported Data Types를 참조하세요. SocketWindowWordCount를 기반으로 POJO와 Tuple 데이터 타입 각각에 대해 두 가지 구현을 보여줍니다.

이 모든 예시에서 연결된 Keyspace example과 Table wordcount가 생성되었다고 가정합니다.

CREATE KEYSPACE IF NOT EXISTS example
    WITH replication = {'class': 'SimpleStrategy', 'replication_factor': '1'};

CREATE TABLE IF NOT EXISTS example.wordcount (
    word text,
    count bigint,
    PRIMARY KEY(word)
);

스트리밍 Tuple 데이터 타입용 Cassandra Sink 예시

Java/Scala Tuple 데이터 타입으로 결과를 Cassandra sink에 저장할 때, 각 레코드를 데이터베이스에 유지하려면 CQL upsert 문을 설정해야 합니다(setQuery('stmt') 사용). upsert 쿼리가 PreparedStatement로 캐시되면 각 Tuple 요소가 문의 매개변수로 변환됩니다.

PreparedStatementBoundStatement에 대한 자세한 내용은 DataStax Java Driver manual을 방문하세요.

Java

// get the execution environment
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// get input data by connecting to the socket
DataStream<String> text = env.socketTextStream(hostname, port, "\n");

// parse the data, group it, window it, and aggregate the counts
DataStream<Tuple2<String, Long>> result = text
        .flatMap(new FlatMapFunction<String, Tuple2<String, Long>>() {
            @Override
            public void flatMap(String value, Collector<Tuple2<String, Long>> out) {
                // normalize and split the line
                String[] words = value.toLowerCase().split("\\s");

                // emit the pairs
                for (String word : words) {
                    //Do not accept empty word, since word is defined as primary key in C* table
                    if (!word.isEmpty()) {
                        out.collect(new Tuple2<String, Long>(word, 1L));
                    }
                }
            }
        })
        .keyBy(value -> value.f0)
        .window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
        .sum(1);

CassandraSink.addSink(result)
        .setQuery("INSERT INTO example.wordcount(word, count) values (?, ?);")
        .setHost("127.0.0.1")
        .build();

Scala

val env: StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment

// get input data by connecting to the socket
val text: DataStream[String] = env.socketTextStream(hostname, port, '\n')

// parse the data, group it, window it, and aggregate the counts
val result: DataStream[(String, Long)] = text
  // split up the lines in pairs (2-tuples) containing: (word,1)
  .flatMap(_.toLowerCase.split("\\s"))
  .filter(_.nonEmpty)
  .map((_, 1L))
  // group by the tuple field "0" and sum up tuple field "1"
  .keyBy(_._1)
  .window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
  .sum(1)

CassandraSink.addSink(result)
  .setQuery("INSERT INTO example.wordcount(word, count) values (?, ?);")
  .setHost("127.0.0.1")
  .build()

result.print().setParallelism(1)

스트리밍 POJO 데이터 타입용 Cassandra Sink 예시

POJO 데이터 타입을 스트리밍하고 동일한 POJO 엔티티를 Cassandra에 저장하는 예시입니다. 이 POJO 구현은 DataStax Java Driver Manual을 따라 클래스에 주석을 달아야 하며, 엔티티의 각 필드는 DataStax Java Driver com.datastax.driver.mapping.Mapper 클래스를 사용해 지정된 테이블의 해당 열에 매핑됩니다.

각 테이블 열의 매핑은 Pojo 클래스의 필드 선언에 배치된 주석을 통해 정의할 수 있습니다. 매핑에 대한 자세한 내용은 Definition of Mapped ClassesCQL Data types에 대한 CQL 문서를 참조하세요.

Java

// get the execution environment
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// get input data by connecting to the socket
DataStream<String> text = env.socketTextStream(hostname, port, "\n");

// parse the data, group it, window it, and aggregate the counts
DataStream<WordCount> result = text
        .flatMap(new FlatMapFunction<String, WordCount>() {
            public void flatMap(String value, Collector<WordCount> out) {
                // normalize and split the line
                String[] words = value.toLowerCase().split("\\s");

                // emit the pairs
                for (String word : words) {
                    if (!word.isEmpty()) {
                        //Do not accept empty word, since word is defined as primary key in C* table
                        out.collect(new WordCount(word, 1L));
                    }
                }
            }
        })
        .keyBy(WordCount::getWord)
        .window(TumblingProcessingTimeWindows.of(Time.seconds(5)))

        .reduce(new ReduceFunction<WordCount>() {
            @Override
            public WordCount reduce(WordCount a, WordCount b) {
                return new WordCount(a.getWord(), a.getCount() + b.getCount());
            }
        });

CassandraSink.addSink(result)
        .setHost("127.0.0.1")
        .setMapperOptions(() -> new Mapper.Option[]{Mapper.Option.saveNullFields(true)})
        .build();


@Table(keyspace = "example", name = "wordcount")
public class WordCount {

    @Column(name = "word")
    private String word = "";

    @Column(name = "count")
    private long count = 0;

    public WordCount() {}

    public WordCount(String word, long count) {
        this.setWord(word);
        this.setCount(count);
    }

    public String getWord() {
        return word;
    }

    public void setWord(String word) {
        this.word = word;
    }

    public long getCount() {
        return count;
    }

    public void setCount(long count) {
        this.count = count;
    }

    @Override
    public String toString() {
        return getWord() + " : " + getCount();
    }
}

더 알아보기 (Learn more)