Mengonsumsi data dari Log Service secara real time dengan SDK mengharuskan Anda mengelola detail implementasi seperti load balancing dan failover antar konsumen. Kelompok konsumen menangani kompleksitas tersebut untuk Anda, sehingga memungkinkan konsumsi data hampir real time, biasanya dalam hitungan detik.
Ikhtisar
Logstore berisi beberapa shard. Kelompok konsumen mengonsumsi data dengan menetapkan shard-shard tersebut ke konsumennya berdasarkan aturan berikut:
Ketika konsumen baru bergabung ke dalam kelompok konsumen, shard akan diseimbangkan ulang di antara semua konsumen untuk memastikan load balancing. Aturan yang sama tetap berlaku.
Konsep utama
|
Term
|
Description
|
|
consumer group
|
Kelompok konsumen terdiri dari beberapa konsumen yang bersama-sama mengonsumsi data dari Logstore yang sama tanpa duplikasi.
Penting
Anda dapat membuat hingga 30 kelompok konsumen untuk setiap Logstore.
|
|
consumer
|
Unit dasar dalam kelompok konsumen yang melakukan konsumsi data sesungguhnya.
Penting
Konsumen dalam kelompok konsumen yang sama harus memiliki nama unik.
|
|
Logstore
|
Unit untuk mengumpulkan, menyimpan, dan mengkueri data. Untuk informasi selengkapnya, lihat Logstore.
|
|
shard
|
Unit yang mengontrol kapasitas baca dan tulis Logstore. Data selalu disimpan dalam sebuah shard. Untuk informasi selengkapnya, lihat shard.
|
|
checkpoint
|
Posisi dalam aliran data yang menandai data terbaru yang telah diproses oleh konsumen. Hal ini memungkinkan konsumen melanjutkan pemrosesan dari titik tersebut setelah restart.
Catatan
Saat Anda mengonsumsi data menggunakan kelompok konsumen, checkpoint secara otomatis disimpan jika program gagal. Setelah program pulih, konsumsi data dapat dilanjutkan dari checkpoint terakhir. Hal ini mencegah konsumsi ganda.
|
Langkah 1: Membuat kelompok konsumen
Catatan
Anda tidak dapat membuat kelompok konsumen langsung di Konsol Web SLS. Anda hanya dapat membuat kelompok konsumen menggunakan SDK, API, atau CLI.
Anda dapat membuat kelompok konsumen menggunakan SDK, API, atau CLI.
SDK
Kode berikut membuat kelompok konsumen:
CreateConsumerGroup.java
import com.aliyun.openservices.log.Client;
import com.aliyun.openservices.log.common.ConsumerGroup;
import com.aliyun.openservices.log.exception.LogException;
public class CreateConsumerGroup {
public static void main(String[] args) throws LogException {
// Contoh ini memperoleh ID AccessKey dan Rahasia AccessKey dari variabel lingkungan.
String accessId = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID");
String accessKey = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET");
// Masukkan nama Proyek.
String projectName = "ali-test-project";
// Masukkan nama Logstore.
String logstoreName = "ali-test-logstore";
// Tetapkan titik akhir untuk Simple Log Service. Contoh ini menggunakan titik akhir Wilayah Tiongkok (Hangzhou). Ganti dengan titik akhir yang sebenarnya.
String host = "https://cn-hangzhou.log.aliyuncs.com";
// Buat client Simple Log Service.
Client client = new Client(host, accessId, accessKey);
try {
// Tetapkan nama kelompok konsumen.
String consumerGroupName = "ali-test-consumergroup2";
System.out.println("ready to create consumergroup");
ConsumerGroup consumerGroup = new ConsumerGroup(consumerGroupName, 300, true);
client.CreateConsumerGroup(projectName, logstoreName, consumerGroup);
System.out.println(String.format("create consumergroup %s success", consumerGroupName));
} catch (LogException e) {
System.out.println("LogException e :" + e.toString());
System.out.println("error code :" + e.GetErrorCode());
System.out.println("error message :" + e.GetErrorMessage());
throw e;
}
}
}
Untuk contoh kode pengelolaan kelompok konsumen, lihat Use the Java SDK to manage consumer groups dan Use Simple Log Service SDK for Python to manage consumer groups.
Langkah 2: Mengonsumsi data log
Cara kerja
Ketika konsumen yang menggunakan SDK kelompok konsumen pertama kali dijalankan, SDK akan membuat kelompok konsumen jika belum ada. Titik mulai checkpoint menentukan posisi konsumsi awal, yang hanya digunakan saat kelompok konsumen pertama kali dibuat. Pada restart berikutnya, konsumen melanjutkan dari checkpoint terakhir yang disimpan oleh server. Contohnya:
-
LogHubConfig.ConsumePosition.BEGIN_CURSOR: Kelompok konsumen mulai mengonsumsi dari log pertama di Logstore.
-
LogHubConfig.ConsumePosition.END_CURSOR: Kelompok konsumen mulai mengonsumsi setelah log terakhir di Logstore.
Contoh
Anda dapat mengonsumsi data dengan kelompok konsumen menggunakan SDK Java, C++, Python, dan Go. Contoh berikut menggunakan Java.
Contoh 1: Konsumsi menggunakan SDK
-
Tambahkan dependensi Maven.
Dalam file pom.xml Anda, tambahkan dependensi berikut:
<dependency>
<groupId>com.google.protobuf</groupId>
<artifactId>protobuf-java</artifactId>
<version>2.5.0</version>
</dependency>
<dependency>
<groupId>com.aliyun.openservices</groupId>
<artifactId>loghub-client-lib</artifactId>
<version>0.6.50</version>
</dependency>
-
Buat kelas untuk menentukan logika pemrosesan log Anda.
SampleLogHubProcessor.java
import com.aliyun.openservices.log.common.FastLog;
import com.aliyun.openservices.log.common.FastLogContent;
import com.aliyun.openservices.log.common.FastLogGroup;
import com.aliyun.openservices.log.common.FastLogTag;
import com.aliyun.openservices.log.common.LogGroupData;
import com.aliyun.openservices.loghub.client.ILogHubCheckPointTracker;
import com.aliyun.openservices.loghub.client.exceptions.LogHubCheckPointException;
import com.aliyun.openservices.loghub.client.interfaces.ILogHubProcessor;
import java.util.List;
public class SampleLogHubProcessor implements ILogHubProcessor {
private int shardId;
// Melacak waktu terakhir checkpoint disimpan.
private long mLastSaveTime = 0;
// Dipanggil sekali saat inisialisasi prosesor.
public void initialize(int shardId) {
this.shardId = shardId;
}
// Logika utama untuk mengonsumsi data. Semua pengecualian harus ditangani dalam metode ini. Jangan melempar pengecualian secara langsung.
public String process(List<LogGroupData> logGroups, ILogHubCheckPointTracker checkPointTracker) {
// Cetak data yang diambil.
for (LogGroupData logGroup : logGroups) {
FastLogGroup fastLogGroup = logGroup.GetFastLogGroup();
System.out.println("Tags");
for (int i = 0; i < fastLogGroup.getLogTagsCount(); ++i) {
FastLogTag logTag = fastLogGroup.getLogTags(i);
System.out.printf("%s : %s\n", logTag.getKey(), logTag.getValue());
}
for (int i = 0; i < fastLogGroup.getLogsCount(); ++i) {
FastLog log = fastLogGroup.getLogs(i);
System.out.println("--------\nLog: " + i + ", time: " + log.getTime() + ", GetContentCount: " + log.getContentsCount());
for (int j = 0; j < log.getContentsCount(); ++j) {
FastLogContent content = log.getContents(j);
System.out.println(content.getKey() + "\t:\t" + content.getValue());
}
}
}
long curTime = System.currentTimeMillis();
// Simpan checkpoint ke server setiap 30 detik. Jika pekerja berhenti secara tak terduga, pekerja baru akan melanjutkan dari checkpoint terakhir, berpotensi memproses ulang sejumlah kecil data.
try {
if (curTime - mLastSaveTime > 30 * 1000) {
// Parameter 'true' segera menyimpan checkpoint ke server. Secara default, checkpoint di memori secara otomatis disimpan ke server setiap 60 detik.
checkPointTracker.saveCheckPoint(true);
mLastSaveTime = curTime;
} else {
// Parameter 'false' menyimpan checkpoint secara lokal. Checkpoint ini akan disimpan ke server oleh mekanisme pembaruan otomatis.
checkPointTracker.saveCheckPoint(false);
}
} catch (LogHubCheckPointException e) {
e.printStackTrace();
}
return null;
}
// Metode ini dipanggil saat pekerja dimatikan. Anda dapat melakukan tugas pembersihan di sini.
public void shutdown(ILogHubCheckPointTracker checkPointTracker) {
// Simpan checkpoint ke server segera.
try {
checkPointTracker.saveCheckPoint(true);
} catch (LogHubCheckPointException e) {
e.printStackTrace();
}
}
}
Untuk contoh kode lainnya, lihat repositori aliyun-log-consumer-java dan Aliyun LOG Go Consumer.
-
Buat factory untuk menghasilkan instance prosesor log Anda.
SampleLogHubProcessorFactory.java
import com.aliyun.openservices.loghub.client.interfaces.ILogHubProcessor;
import com.aliyun.openservices.loghub.client.interfaces.ILogHubProcessorFactory;
class SampleLogHubProcessorFactory implements ILogHubProcessorFactory {
public ILogHubProcessor generatorProcessor() {
// Hasilkan instance konsumen. Catatan: Setiap pemanggilan metode generatorProcessor harus mengembalikan objek SampleLogHubProcessor yang baru.
return new SampleLogHubProcessor();
}
}
-
Buat kelas utama untuk mengonfigurasi dan menjalankan thread pekerja.
Main.java
import com.aliyun.openservices.loghub.client.ClientWorker;
import com.aliyun.openservices.loghub.client.config.LogHubConfig;
import com.aliyun.openservices.loghub.client.exceptions.LogHubClientWorkerException;
public class Main {
// Titik akhir Simple Log Service. Ganti dengan titik akhir Anda yang sebenarnya.
private static String Endpoint = "cn-hangzhou.log.aliyuncs.com";
// Nama proyek. Ganti dengan nama proyek yang sudah ada.
private static String Project = "ali-test-project";
// Nama Logstore. Ganti dengan nama Logstore yang sudah ada.
private static String Logstore = "ali-test-logstore";
// Anda tidak perlu membuat kelompok konsumen terlebih dahulu; program akan membuatnya secara otomatis saat pertama kali dijalankan.
private static String ConsumerGroup = "ali-test-consumergroup2";
// Contoh ini mengambil ID AccessKey dan Rahasia AccessKey dari variabel lingkungan.
private static String AccessKeyId= System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID");
private static String AccessKeySecret = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET");
public static void main(String[] args) throws LogHubClientWorkerException, InterruptedException {
// "consumer_1" adalah nama konsumen, yang harus unik dalam satu kelompok konsumen. Untuk menyeimbangkan konsumsi di beberapa mesin, Anda dapat menggunakan pengenal unik seperti alamat IP mesin sebagai nama konsumen.
// maxFetchLogGroupSize menetapkan jumlah maksimum LogGroup yang diambil per permintaan. Nilai default biasanya sudah cukup. Anda dapat menyesuaikannya dengan menggunakan config.setMaxFetchLogGroupSize(100). Rentang nilai yang valid adalah (0, 1000].
LogHubConfig config = new LogHubConfig(ConsumerGroup, "consumer_1", Endpoint, Project, Logstore, AccessKeyId, AccessKeySecret, LogHubConfig.ConsumePosition.BEGIN_CURSOR,1000);
ClientWorker worker = new ClientWorker(new SampleLogHubProcessorFactory(), config);
Thread thread = new Thread(worker);
// Setelah thread dimulai, ClientWorker berjalan secara otomatis. ClientWorker mengimplementasikan antarmuka Runnable.
thread.start();
Thread.sleep(60 * 60 * 1000);
// Panggil metode shutdown() pekerja untuk menghentikan instance konsumen dan thread terkaitnya.
worker.shutdown();
// ClientWorker membuat tugas asinkron. Setelah shutdown(), tunggu tugas-tugas tersebut selesai secara elegan. Disarankan untuk memberi jeda selama 30 detik.
Thread.sleep(30 * 1000);
}
}
-
Jalankan Main.java.
Berikut adalah contoh output dari konsumsi log NGINX:
: GET
request_uri : /request/path-3/file-7
status : 200
body_bytes_sent : 3820
host : www.example.com
request_time : 43
request_length : 1987
http_user_agent : Mozilla/5.0 (Windows NT 6.1) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/41.0.2228.0 Safari/537.36
http_referer : www.example.com
http_x_forwarded_for : 192.168.10.196
upstream_response_time : 0.02
--------
Log: 158, time: 1635629778, GetContentCount: 14
......
category : null
source : 127.0.0.1
topic : nginx_access_log
machineUUID : null
Tags
__receive_time__ : 1635629815
--------
Log: 0, time: 1635629788, GetContentCount: 14
......
category : null
source : 127.0.0.1
topic : nginx_access_log
machineUUID : null
Tags
__receive_time__ : 1635629877
--------
......
Contoh 2: Konsumsi menggunakan SDK dengan SPL
-
Tambahkan dependensi Maven.
Dalam file pom.xml Anda, tambahkan dependensi berikut:
<dependency>
<groupId>com.google.protobuf</groupId>
<artifactId>protobuf-java</artifactId>
<version>2.5.0</version>
</dependency>
<dependency>
<groupId>com.aliyun.openservices</groupId>
<artifactId>loghub-client-lib</artifactId>
<version>0.6.50</version>
</dependency>
-
Buat kelas untuk menentukan logika pemrosesan log Anda.
SPLLogHubProcessor.java
import com.aliyun.openservices.log.common.FastLog;
import com.aliyun.openservices.log.common.FastLogContent;
import com.aliyun.openservices.log.common.FastLogGroup;
import com.aliyun.openservices.log.common.FastLogTag;
import com.aliyun.openservices.log.common.LogGroupData;
import com.aliyun.openservices.loghub.client.ILogHubCheckPointTracker;
import com.aliyun.openservices.loghub.client.exceptions.LogHubCheckPointException;
import com.aliyun.openservices.loghub.client.interfaces.ILogHubProcessor;
import java.util.List;
public class SPLLogHubProcessor implements ILogHubProcessor {
private int shardId;
// Melacak waktu terakhir checkpoint disimpan.
private long mLastSaveTime = 0;
// Dipanggil sekali saat inisialisasi prosesor.
public void initialize(int shardId) {
this.shardId = shardId;
}
// Logika utama untuk mengonsumsi data. Semua pengecualian harus ditangani dalam metode ini. Jangan melempar pengecualian secara langsung.
public String process(List<LogGroupData> logGroups, ILogHubCheckPointTracker checkPointTracker) {
// Cetak data yang diambil.
for (LogGroupData logGroup : logGroups) {
FastLogGroup fastLogGroup = logGroup.GetFastLogGroup();
System.out.println("Tags");
for (int i = 0; i < fastLogGroup.getLogTagsCount(); ++i) {
FastLogTag logTag = fastLogGroup.getLogTags(i);
System.out.printf("%s : %s\n", logTag.getKey(), logTag.getValue());
}
for (int i = 0; i < fastLogGroup.getLogsCount(); ++i) {
FastLog log = fastLogGroup.getLogs(i);
System.out.println("--------\nLog: " + i + ", time: " + log.getTime() + ", GetContentCount: " + log.getContentsCount());
for (int j = 0; j < log.getContentsCount(); ++j) {
FastLogContent content = log.getContents(j);
System.out.println(content.getKey() + "\t:\t" + content.getValue());
}
}
}
long curTime = System.currentTimeMillis();
// Simpan checkpoint ke server setiap 30 detik. Jika pekerja berhenti secara tak terduga, pekerja baru akan melanjutkan dari checkpoint terakhir, berpotensi memproses ulang sejumlah kecil data.
try {
if (curTime - mLastSaveTime > 30 * 1000) {
// Parameter 'true' segera menyimpan checkpoint ke server. Secara default, checkpoint di memori secara otomatis disimpan ke server setiap 60 detik.
checkPointTracker.saveCheckPoint(true);
mLastSaveTime = curTime;
} else {
// Parameter 'false' menyimpan checkpoint secara lokal. Checkpoint ini akan disimpan ke server oleh mekanisme pembaruan otomatis.
checkPointTracker.saveCheckPoint(false);
}
} catch (LogHubCheckPointException e) {
e.printStackTrace();
}
return null;
}
// Metode ini dipanggil saat pekerja dimatikan. Anda dapat melakukan tugas pembersihan di sini.
public void shutdown(ILogHubCheckPointTracker checkPointTracker) {
// Simpan checkpoint ke server segera.
try {
checkPointTracker.saveCheckPoint(true);
} catch (LogHubCheckPointException e) {
e.printStackTrace();
}
}
}
-
Buat factory untuk menghasilkan instance prosesor log Anda.
SPLLogHubProcessorFactory.java
import com.aliyun.openservices.loghub.client.interfaces.ILogHubProcessor;
import com.aliyun.openservices.loghub.client.interfaces.ILogHubProcessorFactory;
class SPLLogHubProcessorFactory implements ILogHubProcessorFactory {
public ILogHubProcessor generatorProcessor() {
// Hasilkan instance konsumen. Catatan: Setiap pemanggilan metode generatorProcessor harus mengembalikan objek SPLLogHubProcessor yang baru.
return new SPLLogHubProcessor();
}
}
-
Buat kelas utama untuk mengonfigurasi dan menjalankan thread pekerja.
Main.java
import com.aliyun.openservices.loghub.client.ClientWorker;
import com.aliyun.openservices.loghub.client.config.LogHubConfig;
import com.aliyun.openservices.loghub.client.exceptions.LogHubClientWorkerException;
public class Main {
// Titik akhir Simple Log Service. Ganti dengan titik akhir Anda yang sebenarnya.
private static String Endpoint = "cn-hangzhou.log.aliyuncs.com";
// Nama proyek. Ganti dengan nama proyek yang sudah ada.
private static String Project = "ali-test-project";
// Nama Logstore. Ganti dengan nama Logstore yang sudah ada.
private static String Logstore = "ali-test-logstore";
// Anda tidak perlu membuat kelompok konsumen terlebih dahulu; program akan membuatnya secara otomatis saat pertama kali dijalankan.
private static String ConsumerGroup = "ali-test-consumergroup2";
// Contoh ini mengambil ID AccessKey dan Rahasia AccessKey dari variabel lingkungan.
private static String AccessKeyId = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID");
private static String AccessKeySecret = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET");
public static void main(String[] args) throws LogHubClientWorkerException, InterruptedException {
// "consumer_1" adalah nama konsumen, yang harus unik dalam satu kelompok konsumen. Untuk menyeimbangkan konsumsi di beberapa mesin, Anda dapat menggunakan pengenal unik seperti alamat IP mesin sebagai nama konsumen.
// maxFetchLogGroupSize menetapkan jumlah maksimum LogGroup yang diambil per permintaan. Nilai default biasanya sudah cukup. Anda dapat menyesuaikannya dengan menggunakan config.setMaxFetchLogGroupSize(100). Rentang nilai yang valid adalah (0, 1000].
LogHubConfig config = new LogHubConfig(ConsumerGroup, "consumer_1", Endpoint, Project, Logstore, AccessKeyId, AccessKeySecret, LogHubConfig.ConsumePosition.BEGIN_CURSOR, 1000);
// Gunakan setQuery untuk menentukan pernyataan SPL guna memfilter log selama konsumsi.
config.setQuery("* | where cast(body_bytes_sent as bigint) > 14000");
ClientWorker worker = new ClientWorker(new SPLLogHubProcessorFactory(), config);
Thread thread = new Thread(worker);
// Setelah thread dimulai, ClientWorker berjalan secara otomatis. ClientWorker mengimplementasikan antarmuka Runnable.
thread.start();
Thread.sleep(60 * 60 * 1000);
// Panggil metode shutdown() pekerja untuk menghentikan instance konsumen dan thread terkaitnya.
worker.shutdown();
// ClientWorker membuat tugas asinkron. Setelah shutdown(), tunggu tugas-tugas tersebut selesai secara elegan. Disarankan untuk memberi jeda selama 30 detik.
Thread.sleep(30 * 1000);
}
}
-
Jalankan Main.java.
Berikut adalah contoh output dari konsumsi log NGINX:
: GET
request_uri : /request/path-3/file-7
status : 200
body_bytes_sent : 3820
host : www.example.com
request_time : 43
request_length : 1987
http_user_agent : Mozilla/5.0 (Windows NT 6.1) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/41.0.2228.0 Safari/537.36
http_referer : www.example.com
http_x_forwarded_for : 192.168.10.196
upstream_response_time : 0.02
--------
Log: 158, time: 1635629778, GetContentCount: 14
......
category : null
source : 127.0.0.1
topic : nginx_access_log
machineUUID : null
Tags
__receive_time__ : 1635629815
--------
Log: 0, time: 1635629788, GetContentCount: 14
......
category : null
source : 127.0.0.1
topic : nginx_access_log
machineUUID : null
Tags
__receive_time__ : 1635629877
--------
......
Langkah 3: Melihat status kelompok konsumen
Anda dapat melihat status kelompok konsumen dengan salah satu metode berikut:
Java SDK
-
Lihat titik pemeriksaan konsumsi untuk setiap shard. Kode berikut memberikan contoh:
ConsumerGroupTest.java
import java.util.List;
import com.aliyun.openservices.log.Client;
import com.aliyun.openservices.log.common.Consts.CursorMode;
import com.aliyun.openservices.log.common.ConsumerGroup;
import com.aliyun.openservices.log.common.ConsumerGroupShardCheckPoint;
import com.aliyun.openservices.log.exception.LogException;
public class ConsumerGroupTest {
static String endpoint = "cn-hangzhou.log.aliyuncs.com";
static String project = "ali-test-project";
static String logstore = "ali-test-logstore";
static String accesskeyId = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID");
static String accesskey = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET");
public static void main(String[] args) throws LogException {
Client client = new Client(endpoint, accesskeyId, accesskey);
// Dapatkan semua kelompok konsumen di Logstore. Daftar kosong jika tidak ada kelompok konsumen.
List<ConsumerGroup> consumerGroups = client.ListConsumerGroup(project, logstore).GetConsumerGroups();
for(ConsumerGroup c: consumerGroups){
// Cetak properti kelompok konsumen: nama, timeout heartbeat, dan status konsumsi terurut.
System.out.println("Name: " + c.getConsumerGroupName());
System.out.println("heartbeat timeout: " + c.getTimeout());
System.out.println("ordered consumption: " + c.isInOrder());
for(ConsumerGroupShardCheckPoint cp: client.GetCheckPoint(project, logstore, c.getConsumerGroupName()).GetCheckPoints()){
System.out.println("shard: " + cp.getShard());
// Waktu berupa bilangan bulat panjang, akurat hingga mikrodetik.
System.out.println("Last checkpoint update time: " + cp.getUpdateTime());
System.out.println("Consumer name: " + cp.getConsumer());
String consumerPrg = "";
if(cp.getCheckPoint().isEmpty())
consumerPrg = "Consumption has not started";
else{
// Stempel waktu UNIX dalam satuan detik. Format output sesuai kebutuhan.
try{
int prg = client.GetPrevCursorTime(project, logstore, cp.getShard(), cp.getCheckPoint()).GetCursorTime();
consumerPrg = "" + prg;
}
catch(LogException e){
if(e.GetErrorCode() == "InvalidCursor")
consumerPrg = "Invalid. The consumption checkpoint is older than the data retention period.";
else{
// internal server error
throw e;
}
}
}
System.out.println("consumption checkpoint: " + consumerPrg);
String endCursor = client.GetCursor(project, logstore, cp.getShard(), CursorMode.END).GetCursor();
int endPrg = 0;
try{
endPrg = client.GetPrevCursorTime(project, logstore, cp.getShard(), endCursor).GetCursorTime();
}
catch(LogException e){
// do nothing
}
// Stempel waktu UNIX dalam satuan detik. Format output sesuai kebutuhan.
System.out.println("Arrival time of the last record: " + endPrg);
}
}
}
}
-
Berikut adalah contoh output:
Name: ali-test-consumergroup2
heartbeat timeout: 60
ordered consumption: false
shard: 0
Last checkpoint update time: 0
Consumer name: consumer_1
consumption checkpoint: Consumption has not started
Arrival time of the last record: 1729583617
shard: 1
Last checkpoint update time: 0
Consumer name: consumer_1
consumption checkpoint: Consumption has not started
Arrival time of the last record: 1729583738
Process finished with exit code 0
Konsol
-
Login ke Konsol Simple Log Service.
Pada bagian Projects, klik proyek yang diinginkan.

-
Pada tab , klik ikon
di sebelah kiri Logstore target, lalu klik ikon
di sebelah kiri Data Consumption.
-
Pada daftar kelompok konsumen, klik kelompok konsumen target.
-
Pada halaman Consumer Group Status, lihat titik pemeriksaan konsumsi untuk setiap shard. Halaman ini menampilkan detail untuk setiap shard, termasuk ID-nya (shard), Last Consumed Time, dan Client yang mengonsumsi. Anda juga dapat menggunakan tombol Refresh dan Reset Checkpoint.
Operasi terkait
-
Otorisasi Pengguna RAM
Untuk menggunakan Pengguna RAM dalam mengelola kelompok konsumen, berikan izin yang diperlukan kepada Pengguna RAM tersebut. Untuk informasi selengkapnya, lihat Create and authorize a RAM user.
Tabel berikut mencantumkan Actions yang diperlukan.
|
Action
|
Description
|
Resource
|
|
log:GetCursorOrData(GetCursor)
|
Mendapatkan cursor berdasarkan waktu tertentu.
|
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}
|
|
log:CreateConsumerGroup(CreateConsumerGroup)
|
Membuat kelompok konsumen di Logstore tertentu.
|
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}/consumergroup/${consumerGroupName}
|
|
log:ListConsumerGroup(ListConsumerGroup)
|
Menampilkan semua kelompok konsumen di Logstore tertentu.
|
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}/consumergroup/*
|
|
log:ConsumerGroupUpdateCheckPoint(UpdateCheckPoint)
|
Memperbarui checkpoint pada shard untuk kelompok konsumen tertentu.
|
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}/consumergroup/${consumerGroupName}
|
|
log:ConsumerGroupHeartBeat(ConsumerGroupHeartBeat)
|
Mengirim heartbeat dari konsumen tertentu ke server.
|
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}/consumergroup/${consumerGroupName}
|
|
log:UpdateConsumerGroup(UpdateConsumerGroup)
|
Memodifikasi properti kelompok konsumen tertentu.
|
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}/consumergroup/${consumerGroupName}
|
|
log:GetConsumerGroupCheckPoint(GetConsumerGroupCheckPoint)
|
Mendapatkan checkpoint satu atau semua shard untuk kelompok konsumen tertentu.
|
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}/consumergroup/${consumerGroupName}
|
Untuk memberikan izin yang tercantum di bawah ini kepada Pengguna RAM, gunakan kebijakan contoh berikut.
-
ID akun Alibaba Cloud: 174649****602745
-
ID wilayah: cn-hangzhou
-
Nama proyek: project-test
-
Nama Logstore: logstore-test
-
Nama kelompok konsumen: consumergroup-test
Contoh kebijakan:
{
"Version": "1",
"Statement": [
{
"Effect": "Allow",
"Action": [
"log:GetCursorOrData"
],
"Resource": "acs:log:cn-hangzhou:174649****602745:project/project-test/logstore/logstore-test"
},
{
"Effect": "Allow",
"Action": [
"log:CreateConsumerGroup",
"log:ListConsumerGroup"
],
"Resource": "acs:log:cn-hangzhou:174649****602745:project/project-test/logstore/logstore-test/consumergroup/*"
},
{
"Effect": "Allow",
"Action": [
"log:ConsumerGroupUpdateCheckPoint",
"log:ConsumerGroupHeartBeat",
"log:UpdateConsumerGroup",
"log:GetConsumerGroupCheckPoint"
],
"Resource": "acs:log:cn-hangzhou:174649****602745:project/project-test/logstore/logstore-test/consumergroup/consumergroup-test"
}
]
}
Sebagai alternatif, Anda dapat menyambungkan kebijakan sistem AliyunLogFullAccess, yang memberikan izin untuk mengelola Simple Log Service, ke Pengguna RAM. Kebijakan ini sudah mencakup semua izin yang diperlukan untuk operasi kelompok konsumen, sehingga Anda tidak perlu membuat kebijakan kustom.
-
Troubleshooting
Untuk mempermudah troubleshooting, konfigurasikan Log4j untuk aplikasi konsumen Anda agar mencatat pengecualian dari kelompok konsumen. Berikut adalah contoh konfigurasi log4j.properties:
log4j.rootLogger = info,stdout
log4j.appender.stdout = org.apache.log4j.ConsoleAppender
log4j.appender.stdout.Target = System.out
log4j.appender.stdout.layout = org.apache.log4j.PatternLayout
log4j.appender.stdout.layout.ConversionPattern = [%-5p] %d{yyyy-MM-dd HH:mm:ss,SSS} method:%l%n%m%n
Setelah Anda mengonfigurasi Log4j, aplikasi konsumen akan menghasilkan informasi pengecualian seperti berikut:
[WARN ] 2018-03-14 12:01:52,747 method:com.aliyun.openservices.loghub.client.LogHubConsumer.sampleLogError(LogHubConsumer.java:159)
com.aliyun.openservices.log.exception.LogException: Invalid loggroup count, (0,1000]
-
Mengonsumsi data dari waktu tertentu
// consumerStartTimeInSeconds menunjukkan waktu mulai konsumsi data.
public LogHubConfig(String consumerGroupName,
String consumerName,
String loghubEndPoint,
String project, String logStore,
String accessId, String accessKey,
int consumerStartTimeInSeconds);
// position adalah enumerasi. LogHubConfig.ConsumePosition.BEGIN_CURSOR memulai konsumsi dari data terlama. LogHubConfig.ConsumePosition.END_CURSOR memulai konsumsi dari data terbaru.
public LogHubConfig(String consumerGroupName,
String consumerName,
String loghubEndPoint,
String project, String logStore,
String accessId, String accessKey,
ConsumePosition position);
Catatan
-
Pilih konstruktor sesuai kebutuhan Anda.
-
Jika checkpoint sudah disimpan di server, konsumsi dilanjutkan dari checkpoint tersebut.
-
Simple Log Service memprioritaskan checkpoint yang tersimpan untuk konsumsi. Jika Anda menentukan waktu mulai, pastikan nilai consumerStartTimeInSeconds berada dalam periode retensi data (TTL). Jika tidak, waktu mulai yang ditentukan akan diabaikan.
-
Reset checkpoint
public static void updateCheckpoint() throws Exception {
Client client = new Client(host, accessId, accessKey);
// Stempel waktu harus berupa stempel waktu UNIX dalam satuan detik. Jika stempel waktu Anda dalam milidetik, bagi dengan 1000.
long timestamp = Timestamp.valueOf("2017-11-15 00:00:00").getTime() / 1000;
ListShardResponse response = client.ListShard(new ListShardRequest(project, logStore));
for (Shard shard : response.GetShards()) {
int shardId = shard.GetShardId();
String cursor = client.GetCursor(project, logStore, shardId, timestamp).GetCursor();
client.UpdateCheckPoint(project, logStore, consumerGroup, shardId, cursor);
}
}
Referensi
-
API
-
SDK
|
Language
|
References
|
|
Java
|
|
|
Python
|
|
-
CLI