Sink data ke Redis dengan cara Squirrel

Mar 27 2023
Apache Flink dan Redis adalah dua alat canggih yang dapat digunakan bersama untuk membangun saluran pemrosesan data real-time yang dapat menangani volume data yang besar. Flink menyediakan platform yang sangat skalabel dan toleran terhadap kesalahan untuk memproses aliran data, sementara Redis menyediakan database dalam memori berkinerja tinggi yang dapat digunakan untuk menyimpan dan meminta data.

Apache Flink dan Redis adalah dua alat canggih yang dapat digunakan bersama untuk membangun saluran pemrosesan data real-time yang dapat menangani volume data yang besar. Flink menyediakan platform yang sangat skalabel dan toleran terhadap kesalahan untuk memproses aliran data, sementara Redis menyediakan database dalam memori berkinerja tinggi yang dapat digunakan untuk menyimpan dan meminta data. Pada artikel ini, kita akan mengeksplorasi bagaimana Flink dapat digunakan untuk memanggil Redis menggunakan fungsi async dan menunjukkan bagaimana ini dapat digunakan untuk mendorong data ke Redis dengan cara yang tidak memblokir.

Kisah Redis

IC: Penjelasan Redis Infografis(https://architecturenotes.co/redis/)

“Redis: Lebih dari Sekedar Cache

Redis adalah penyimpanan struktur data dalam memori NoSQL yang kuat yang telah menjadi alat bantu bagi pengembang. Meskipun sering dianggap hanya sebagai cache, Redis lebih dari itu. Itu bisa berfungsi sebagai database, perantara pesan, dan cache semuanya dalam satu.

Salah satu kekuatan Redis adalah keserbagunaannya. Ini mendukung berbagai tipe data, termasuk Strings, Lists, Sets, Sorted Sets, Hashes, Streams, HyperLogLogs, dan Bitmaps. Redis juga menawarkan indeks geospasial dan kueri radius, menjadikannya alat yang berharga untuk aplikasi berbasis lokasi.

Fitur Redis melampaui model datanya. Ini memiliki replikasi bawaan, skrip Lua, dan transaksi, dan dapat secara otomatis mempartisi data dengan Redis Cluster. Selain itu, Redis menyediakan ketersediaan tinggi melalui Redis Sentinel.

Catatan: Pada artikel ini, kami akan lebih fokus pada Redis Cluster Mode

IC: Mode Klaster Redis (https://architecturenotes.co/redis/)

Redis Cluster menggunakan sharding algoritmik dengan Hashslots untuk menentukan shard mana yang memegang kunci tertentu dan menyederhanakan penambahan instance baru. Sementara itu, menggunakan Gossiping untuk menentukan kesehatan cluster, dan jika node primer tidak responsif, node sekunder dapat dipromosikan untuk menjaga cluster tetap sehat. Sangat penting untuk memiliki jumlah node primer ganjil dan dua replika untuk penyiapan yang kuat untuk menghindari fenomena otak terbagi (Di mana cluster tidak dapat memutuskan siapa yang akan dipromosikan dan berakhir dengan keputusan terpisah)

Untuk berbicara dengan Redis Cluster kita akan menggunakan klien Redis Async Java.

Kisah Flink

IC: Flink Tingkat Tinggi (https://flink.apache.org/)

Apache Flink adalah open-source, kerangka pemrosesan aliran terpadu dan pemrosesan batch yang dirancang untuk menangani pemrosesan data real-time, throughput tinggi, dan toleran terhadap kesalahan. Itu dibangun di atas kerangka kerja Apache Gelly dan dirancang untuk mendukung pemrosesan peristiwa yang kompleks dan perhitungan stateful pada Bounded dan Unbounded Streams. Apa yang membuatnya cepat adalah Memanfaatkan Kinerja Dalam Memori dan secara asinkron memeriksa keadaan lokal.

Pahlawan Cerita

IC: Flink 1.16 Rilis Docs

Interaksi asinkron dengan database adalah pengubah permainan untuk aplikasi pemrosesan aliran. Dengan pendekatan ini, instance fungsi tunggal dapat menangani beberapa permintaan sekaligus, memungkinkan respons bersamaan dan peningkatan throughput yang signifikan. Dengan tumpang tindih waktu tunggu dengan permintaan dan tanggapan lain, alur pemrosesan menjadi jauh lebih efisien.

Kami akan mengambil Contoh Data E-niaga untuk Menghitung Jumlah penjualan untuk Setiap kategori dalam Jendela Geser 24 Jam dengan slide 30 Detik dan Menenggelamkannya ke Redis untuk Pencarian yang lebih cepat untuk layanan hilir.

Contoh Kumpulan Data



Category, TimeStamp
Electronics,1679832334
Furniture,1679832336
Fashion,1679832378
Food,16798323536


package Aysnc_kafka_redis;

import AsyncIO.RedisSink;
import akka.japi.tuple.Tuple3;
import deserializer.Ecommdeserialize;
import model.Ecomm;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.streaming.api.datastream.AsyncDataStream;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.timestamps.BoundedOutOfOrdernessTimestampExtractor;
import org.apache.flink.streaming.api.functions.windowing.WindowFunction;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
import org.apache.flink.util.Collector;
import java.util.concurrent.TimeUnit;

public class FlinkAsyncRedis {

    public static void main(String[] args) throws Exception {


        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        Ecommdeserialize jsonde = new Ecommdeserialize(); 

        KafkaSource<Ecomm> source = KafkaSource.<Ecomm>builder()
                .setTopics("{dummytopic}")
                .setBootstrapServers("{dummybootstrap}")
                .setGroupId("test_flink")
                .setStartingOffsets(OffsetsInitializer.earliest())
                .setValueOnlyDeserializer(jsonde)
                .build();


        DataStream<Ecomm> orderData = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source");


        orderData.assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractor<Ecomm>(Time.seconds(10)) {
            @Override
            public long extractTimestamp(Ecomm element) {
                return element.getEventTimestamp(); // extract watermark column from stream
            }
        });

        SingleOutputStreamOperator<Tuple3<String, Long, Long>> aggregatedData = orderData.keyBy(Ecomm::getCategory)
                .window(SlidingEventTimeWindows.of(Time.hours(24),Time.seconds(30)))
                .apply((WindowFunction<Ecomm, Tuple3<String, Long, Long>, String, TimeWindow>) (key, window, input, out) -> {
                    long count = 0;
                    for (Ecomm event : input) {
                        count++; // increment the count for each event in the window
                    }
                    out.collect(new Tuple3<>(key, window.getEnd(), count)); // output the category, window end time, and count
                });


        // calling async I/0 operator to sink data to redis in UnOrdered way
        SingleOutputStreamOperator<String> sinkResults = AsyncDataStream.unorderedWait(aggregatedData,new RedisSink(
                "{redisClusterUrl}"),
                1000, // the timeout defines how long an asynchronous operation take before it is finally considered failed
                TimeUnit.MILLISECONDS,
                 100); //capacity This parameter defines how many asynchronous requests may be in progress at the same time.

        sinkResults.print(); // print out the redis set response stored in the future for every key

        env.execute("RedisAsyncSink"); // you will be able to see your job running on cluster by this name


    }

}

package AsyncIO;

import akka.japi.tuple.Tuple3;
import io.lettuce.core.RedisFuture;
import io.lettuce.core.cluster.RedisClusterClient;
import io.lettuce.core.cluster.api.StatefulRedisClusterConnection;
import io.lettuce.core.cluster.api.async.RedisAdvancedClusterAsyncCommands;
import lombok.AllArgsConstructor;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.async.ResultFuture;
import org.apache.flink.streaming.api.functions.async.RichAsyncFunction;
import scala.collection.immutable.List;

import java.util.ArrayList;
import java.util.Collections;

@AllArgsConstructor
public class RedisSink extends RichAsyncFunction<Tuple3<String, Long, Long>, String> {

    String redisUrl;

    public RedisSink(String redisUrl){
        this.redisUrl=redisUrl;
    }

    private transient RedisClusterClient client = null;
    private transient StatefulRedisClusterConnection<String, String> clusterConnection = null;
    private transient RedisAdvancedClusterAsyncCommands<String, String> asyncCall = null;


    // method executes any operator-specific initialization 
    @Override
    public void open(Configuration parameters) {
        if (client == null ) {
            client = RedisClusterClient.create(redisUrl);
        }
        if (clusterConnection == null) {
            clusterConnection = client.connect();
        }
        if (asyncCall == null) {
            asyncCall  = clusterConnection.async();
        }
    }

    // core logic to set key in redis using async connection and return result of the call via ResultFuture
    @Override
    public void asyncInvoke(Tuple3<String, Long, Long> stream, ResultFuture<String> resultFuture) {

        String productKey = stream.t1();
        System.out.println("RedisKey:" + productKey); //for logging
        String count = stream.t3().toString();
        System.out.println("Redisvalue:" +  count); //for logging
        RedisFuture<String> setResult = asyncCall.set(productKey,count);

        setResult.whenComplete((result, throwable) -> {if(throwable!=null){
            System.out.println("Callback from redis failed:" + throwable);
            resultFuture.complete(new ArrayList<>());
        }
        else{
            resultFuture.complete(new ArrayList(Collections.singleton(result)));
        }});
    }
    
     // method closes what was opened during initialization to free any resources 
    //  held by the operator (e.g. open network connections, io streams)
    @Override
    public void close() throws Exception {
        client.close();
    }


}

  • Data yang dialirkan ke Redis dapat Anda gunakan oleh model ilmu Data untuk mencari dan menghasilkan lebih banyak produk untuk kategori yang sering dijual selama musim obral.
  • Ini dapat digunakan untuk menampilkan grafik dan angka sebagai Stats of the Sale di halaman web, untuk membuat dorongan di antara pengguna untuk pembelian yang agresif.
  • Flink menyediakan platform yang sangat skalabel dan toleran terhadap kesalahan untuk memproses aliran data, sementara Redis menyediakan database dalam memori berkinerja tinggi yang dapat digunakan untuk menyimpan dan meminta data.
  • Pemrograman asinkron dapat digunakan untuk meningkatkan kinerja pipeline pemrosesan data dengan mengizinkan panggilan non-pemblokiran ke sistem eksternal seperti Redis.
  • Kombinasi keduanya dapat membantu menghadirkan budaya keputusan data waktu nyata.

https://architecturenotes.co/redis/.

https://www.baeldung.com/java-redis-lettuce

https://nightlies.apache.org/flink/flink-docs-release-1.16/docs/dev/datastream/operators/asyncio/