/
happy_goose
/
KafkaToPostgres
Обзор
Документация
Войти
/
happy_goose
/
KafkaToPostgres
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
src/main/java/sberdata/KafkaToPostgres.java
128 строк
6 KB
olivan139
Add kafka parallel read and postgres write
16 янв 2025, 17:30
16 янв 2025, 17:30
ff339ad
Код
Авторство
О чём код?
package sberdata; import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.api.common.serialization.SimpleStringSchema; 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.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.sink.SinkFunction; import org.apache.flink.streaming.api.functions.KeyedProcessFunction; import org.apache.flink.util.Collector; import java.sql.Connection; import java.sql.DriverManager; import java.sql.PreparedStatement; public class KafkaToPostgres { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); KafkaSource<String> kafkaSource = KafkaSource.<String>builder() .setBootstrapServers("localhost:9092") .setTopics("bm512orderline3") .setStartingOffsets(OffsetsInitializer.earliest()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); // Read data from Kafka DataStream<String> kafkaStream = env.fromSource( kafkaSource, WatermarkStrategy.noWatermarks(), "Kafka Source" ); // .setParallelism(4) kafkaStream .map(value -> { String[] fields = value.split(","); return new CustomerData( Integer.parseInt(fields[0]), // wid Integer.parseInt(fields[1]), // did Integer.parseInt(fields[2]), // cid Integer.parseInt(fields[3]), // orderid fields[4], // lastname fields[5], // firstname fields[6], // city fields[7], // state fields[8], // zip Double.parseDouble(fields[9]) // amount ); }) .keyBy(data -> data.getCompositeKey()) .process(new UpsertProcessFunction()); env.execute("Kafka to PostgreSQL Upsert Job"); } public static class CustomerData { public int wid; public int did; public int cid; public int orderid; public String lastname; public String firstname; public String city; public String state; public String zip; public double amount; public CustomerData(int wid, int did, int cid, int orderid, String lastname, String firstname, String city, String state, String zip, double amount) { this.wid = wid; this.did = did; this.cid = cid; this.orderid = orderid; this.lastname = lastname; this.firstname = firstname; this.city = city; this.state = state; this.zip = zip; this.amount = amount; } public String getCompositeKey() { return wid + "|" + did + "|" + cid + "|" + orderid; } } public static class UpsertProcessFunction extends KeyedProcessFunction<String, CustomerData, Void> { private static final String DB_URL = "jdbc:postgresql://localhost:5432/postgres"; private static final String DB_USER = "postgres"; private static final String DB_PASSWORD = "password"; @Override public void processElement(CustomerData value, Context ctx, Collector<Void> out) throws Exception { Class.forName("org.postgresql.Driver"); try (Connection connection = DriverManager.getConnection(DB_URL, DB_USER, DB_PASSWORD)) { String upsertQuery = "INSERT INTO customer_dm_3b (c_w_id, c_d_id, c_c_id, c_o_id, c_last, c_first, c_city, c_state, c_zip, c_order_line_cnt, c_order_amount_sum, c_order_amount_min, c_order_amount_max, c_order_amount_avg, c_order_amount_last) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 1, ?, ?, ?, ?, ?) " + "ON CONFLICT (c_w_id, c_d_id, c_c_id, c_o_id) DO UPDATE SET " + "c_order_amount_sum = customer_dm_3b.c_order_amount_sum + EXCLUDED.c_order_amount_last, " + "c_order_amount_min = LEAST(customer_dm_3b.c_order_amount_min, EXCLUDED.c_order_amount_last), " + "c_order_amount_max = GREATEST(customer_dm_3b.c_order_amount_max, EXCLUDED.c_order_amount_last), " + "c_order_amount_avg = (customer_dm_3b.c_order_amount_sum + EXCLUDED.c_order_amount_last) / (customer_dm_3b.c_order_line_cnt + 1), " + "c_order_amount_last = EXCLUDED.c_order_amount_last, " + "c_order_line_cnt = customer_dm_3b.c_order_line_cnt + 1;"; try (PreparedStatement statement = connection.prepareStatement(upsertQuery)) { statement.setInt(1, value.wid); statement.setInt(2, value.did); statement.setInt(3, value.cid); statement.setInt(4, value.orderid); statement.setString(5, value.lastname); statement.setString(6, value.firstname); statement.setString(7, value.city); statement.setString(8, value.state); statement.setString(9, value.zip); statement.setDouble(10, value.amount); // c_order_amount_sum statement.setDouble(11, value.amount); // c_order_amount_min statement.setDouble(12, value.amount); // c_order_amount_max statement.setDouble(13, value.amount); // c_order_amount_avg statement.setDouble(14, value.amount); // c_order_amount_last statement.executeUpdate(); } } } } }