Hubungkan Apache Flink ke LindormTable dan gunakan tabel Lindorm sebagai tabel dimensi atau tabel hasil dalam pekerjaan Flink Anda. Anda dapat mengakses tabel-tabel ini menggunakan Flink SQL atau Flink DataStream.
Informasi latar belakang
Anda dapat menggunakan Lindorm LindormTable sebagai tabel dimensi atau tabel hasil di Flink, serta mengakses LindormTable melalui Flink SQL atau Flink DataStream.
Pilih metode koneksi
Di LindormTable, tabel dibagi menjadi dua jenis berdasarkan cara pembuatannya: tabel HBase (tabel yang dibuat dan ditulis menggunakan API HBase) dan tabel SQL (tabel yang dibuat dan ditulis menggunakan Lindorm SQL). Mengakses kedua jenis tabel ini memerlukan antarmuka yang berbeda. Oleh karena itu, sebelum membuat task Flink, Anda perlu menentukan jenis connector, lalu menentukan alamat koneksi LindormTable berdasarkan jenis connector tersebut.
Tentukan jenis tabel yang akan diakses
Di Lindorm, Anda dapat menggunakan Lindorm SQL untuk menentukan jenis LindormTable yang perlu diakses oleh task Flink yang sedang dikembangkan:
Gunakan lindorm-cli (atau langsung melalui Lindorm Insight atau DMS) untuk terhubung ke mesin tabel lebar Lindorm.
Jalankan pernyataan SQL berikut untuk memeriksa atribut
IS_HBASE_LIKEdari tabel tersebut.Jika nilai atributnya TRUE, tabel tersebut adalah tabel HBase.
Jika nilai atributnya FALSE, tabel tersebut adalah tabel SQL.
SHOW TABLE VARIABLES FROM table_name LIKE 'IS_HBASE_LIKE';
Untuk informasi tentang cara menggunakan
lindorm-cliuntuk terhubung ke mesin tabel lebar, lihat Hubungkan ke dan gunakan mesin tabel lebar dengan Lindorm-cli.Untuk sintaksis lengkap
SHOW TABLE VARIABLE, lihat SHOW VARIABLES.
Pilih jenis connector
Berdasarkan bentuk produk Flink yang Anda pilih, tentukan connector mana yang akan digunakan untuk mengakses LindormTable.
Jenis tabel | Community Flink | Realtime Compute for Apache Flink |
HBase table | (Didukung sebagai tabel dimensi dan tabel hasil) | Cloud-native multi-model database Lindorm connector (Didukung sebagai tabel dimensi dan tabel hasil) |
SQL table | (Didukung sebagai tabel hasil) | Cloud-native multi-model database Lindorm connector (Didukung sebagai tabel dimensi dan tabel hasil) |
Realtime Compute for Apache Flink adalah layanan Flink terkelola di Alibaba Cloud. Untuk informasi selengkapnya, lihat Realtime Compute for Apache Flink. Perhatikan bahwa kluster Flink yang dibangun menggunakan Flink open source di ECS Alibaba Cloud tetap diklasifikasikan sebagai Community Flink dalam tabel di atas.
Dapatkan informasi koneksi LindormTable
Skenario 1: Gunakan open source HBase connector atau cloud-native multi-model database Lindorm connector
Dalam skenario ini, alamat koneksi harus berupa alamat akses HBase Java API (VPC) dari mesin tabel lebar.
Di halaman detail instans Lindorm, klik Database connections pada menu sebelah kiri, lalu pilih tab Wide Table Engine. Di area Connect through HBase-compatible address, dapatkan alamat VPC dalam formatld-<instance-ID>-proxy-lindorm.lindorm.rds.aliyuncs.com:30020. Di area Connect through MySQL-compatible address, dapatkan alamat kompatibel MySQL dalam formatld-<instance-ID>-proxy-sql-lindorm.lindorm.rds.aliyuncs.com:33060.Skenario 2: Gunakan konektor JDBC
Dalam skenario ini, alamat koneksi harus berupa alamat kompatibel MySQL (VPC) dari mesin tabel wide.
Pada halaman detail instans Lindorm, klik Database connections di panel navigasi sebelah kiri, lalu pilih tab Wide Table Engine. Di bagian Connect through HBase-compatible address, Anda dapat melihat alamat VPC untuk HBase Java API (dalam formatld-<instance-ID>-proxy-lindorm.rds.aliyuncs.com:30020).
Jika task Flink menggunakan pengguna Lindorm yang baru dibuat untuk mengakses LindormTable, pastikan pengguna tersebut memiliki izin baca dan tulis pada tabel Flink. Untuk informasi tentang cara memberikan izin, lihat Berikan izin kepada pengguna tertentu.
Untuk penjelasan rinci tentang berbagai alamat koneksi tabel lebar Lindorm, lihat Lihat alamat koneksi mesin tabel lebar.
Metode bagi pekerjaan Flink untuk mengakses LindormTable
Anda dapat mengembangkan dengan pendekatan umum dari framework komputasi real-time yang dipilih. Untuk mengakses LindormTable dalam suatu pekerjaan, rujuk dokumen berikut berdasarkan connector yang telah Anda pilih agar pekerjaan komputasi Anda dapat mengakses LindormTable.
Prasyarat
Kembangkan pekerjaan yang menggunakan Community Flink untuk mengakses LindormTable
Untuk menggunakan open source HBase connector guna mengakses LindormTable, pastikan mesin tabel lebar berada pada versi 2.4.3 atau lebih baru.
Untuk menggunakan open source JDBC connector guna mengakses LindormTable, pastikan mesin tabel lebar berada pada versi 2.6.5.2 atau lebih baru, dan Anda telah Mengaktifkan fitur kompatibilitas MySQL.
Untuk informasi tentang cara melihat atau meningkatkan versi saat ini, lihat Catatan rilis LindormTable dan Tingkatkan versi mesin minor instans Lindorm.
Jika Anda menggunakan Realtime Compute for Apache Flink untuk mengembangkan pekerjaan yang mengakses LindormTable, tidak ada batasan versi pada mesin tabel lebar.
Pastikan lingkungan tempat kluster Flink berada memiliki konektivitas jaringan dengan instans Lindorm, dan alamat IP klien telah ditambahkan ke daftar putih Lindorm. Untuk informasi tentang cara menambahkan alamat IP ke daftar putih, lihat Konfigurasikan daftar putih.
Pengembangan pekerjaan dengan Community Flink
Open source HBase connector
Saat Anda menggunakan Community Flink untuk mengembangkan pekerjaan yang mengakses tabel HBase, jika ingin mengakses tabel melalui Internet, atau jika instans Lindorm target adalah instans single-node Lindorm, Anda harus meningkatkan SDK dan mengubah konfigurasi sebelum melakukan operasi selanjutnya.
Untuk detailnya, lihat Langkah 1 di Hubungkan ke dan gunakan LindormTable menggunakan HBase Java API.
Untuk informasi tentang cara menggunakan open source HBase connector guna membuat tabel dimensi dan tabel hasil, lihat Dokumentasi HBase connector.
Open source JDBC connector
Saat Anda menggunakan open source JDBC connector untuk mengakses LindormTable, saat ini hanya didukung penggunaan LindormTable sebagai tabel hasil. Untuk penggunaan secara umum, rujuk dokumentasi resmi JDBC connector. Namun, ada beberapa poin yang perlu diperhatikan secara khusus, sebagai berikut:
Kebutuhan dependensi
Dibandingkan dengan rentang versi paket dependensi yang relatif luas dalam dokumentasi resmi, versi paket dependensi yang digunakan oleh JDBC connector untuk mengakses LindormTable saat ini dibatasi pada daftar berikut:flink-connector-jdbc-core-4.0.0-2.0.jarflink-connector-jdbc-mysql-4.0.0-2.0.jarmysql-connector-j-8.3.0.jar
Paket dependensi driver JDBC MySQL dapat diunduh dari komunitas.
Parameter JDBC connector
Karena penggunaan sebagai tabel sumber dan tabel dimensi tidak didukung, parameter yang terkait dengan fungsionalitas tabel sumber dan tabel dimensi tidak didukung (seperti parameter dengan awalan scan sepertiscan.fetch-size, dan parameter dengan awalanlookupsepertilookup.cache).
Rekomendasi untuk beberapa parameter JDBC connector:url: Kami menyarankan Anda mengikuti Gunakan Java JDBC APIs untuk mengembangkan aplikasi dalam konfigurasi.
username: Gunakan username yang dibuat di instans Lindorm.
password: Kata sandi dari username tersebut.
connector, table-name: Ikuti rekomendasi komunitas JDBC connector.
parameter sink: Sesuaikan secara detail berdasarkan kondisi aktual pekerjaan.
Pemetaan tipe data
Pemetaan antara tipe data Flink dan tipe data Lindorm umumnya dapat mengikuti pemetaan tipe data MySQL (rujuk bagian Pemetaan tipe data pada JDBC connector). Namun, beberapa tipe data di Lindorm tidak selaras dengan MySQL. Sebagai contoh, tipe MySQL berikut yang diklaim didukung dalam JDBC connector tidak didukung oleh Lindorm:Tipe MEDIUMINT
Tipe DATETIME
Tipe UNSIGNED selain BIGINT UNSIGNED
Contoh berikut menggunakan Flink SQL untuk mendefinisikan pekerjaan yang mengakses LindormTable melalui JDBC connector. Dalam contoh ini, asumsikan sebuah tabel bernama testflink telah didefinisikan di mesin tabel lebar Lindorm.
# Buat tabel Flink dan mulai pekerjaan
CREATE TABLE source_table(
c1 INT,
c2 STRING
) WITH (
'connector' = 'datagen',
'rows-per-second' = '2',
'fields.c2.length' = '5',
'fields.c1.min' = '1',
'fields.c1.max' = '100'
);
CREATE TABLE sink_table(
c1 INT,
c2 STRING
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://ld-xxxxx-proxy-lindorm.lindorm.rds.aliyuncs.com:33060/default?sslMode=disabled&allowPublicKeyRetrieval=true&useServerPrepStmts=true&useLocalSessionState=true&rewriteBatchedStatements=true&cachePrepStmts=true&prepStmtCacheSize=300&prepStmtCacheSqlLimit=50000000',
'username' = 'root',
'password' = 'root',
'table-name' = 'testflink'
);
INSERT INTO sink_table SELECT * FROM source_table;Pengembangan pekerjaan dengan Realtime Compute for Apache Flink
Cloud-native multi-model database Lindorm connector
Dalam pengembangan pekerjaan Realtime Compute for Apache Flink, Anda dapat mengembangkan pekerjaan yang mengakses mesin tabel lebar Lindorm menggunakan Flink SQL. Untuk detail pengembangan pekerjaan Realtime Compute for Apache Flink, lihat Ikhtisar pengembangan pekerjaan.
Untuk informasi tentang cara menggunakan connector guna membuat tabel dimensi dan tabel hasil, lihat Dokumentasi cloud-native multi-model database Lindorm connector.