跳到主要內容

AWS Kinesis Note

Sample

Ec2 Instance(Consumer)
|
+- Worker +
|- Processor <-- shard (shard 的 seqId, partictionId 會被記錄到 DynamoDB)
|- Processor <-- shard (shard 的 seqId, partictionId 會被記錄到 DynamoDB)
|- Processor <-- shard (shard 的 seqId, partictionId 會被記錄到 DynamoDB)
|- Processor <-- shard (shard 的 seqId, partictionId 會被記錄到 DynamoDB)

Ec2 Instance(Consumer)
|
+- Worker +
|- Processor <-- shard (shard 的 seqId, partictionId 會被記錄到 DynamoDB)
|- Processor <-- shard (shard 的 seqId, partictionId 會被記錄到 DynamoDB)
|- Processor <-- shard (shard 的 seqId, partictionId 會被記錄到 DynamoDB)
|- Processor <-- shard (shard 的 seqId, partictionId 會被記錄到 DynamoDB)

Concept

每一台 Ec2 可視為一個 Consumer,每個 Consumer 會有個 Worker 機制把 shard Stream 分配給 Processor,當 shard 的數量調大的時候,KCL(Kinesis Clinet Libary) 裡面的 IRecordProcessorFactory 就會自動建立對應, shard 的 processor。

當有新的 Consumer 被建立時,AWS 會自動做 loadbalance,去平衡每個 worker 要負責的 shards,這時候 Processor 裡面的 shutdown 會被執行。

切記 Consumer 不應該開的比 shard 的數量還大。

留言

這個網誌中的熱門文章

Parse URI query string to Key Value

Parse URI query String to Map 做 urlDecode 處理 沒有任何 query String 回傳 Empty Map 確保只處理 key-value 結構的 query String package com.example.util; import lombok.extern.slf4j.Slf4j; import org.apache.http.client.utils.URIBuilder; import java.io.UnsupportedEncodingException; import java.net.URI; import java.net.URISyntaxException; import java.net.URLDecoder; import java.util.LinkedHashMap; import java.util.Map; import java.util.Objects; /** * Created by jerry on 2017/12/28. * * @author jerry */ @Slf4j public class UriUtil { private UriUtil() { } public static Map splitQuery(final String uri) { Map queryPairs = new LinkedHashMap (); try { final URI uri = new URIBuilder(uri).build(); final String rawQuery = uri.getRawQuery(); log.info("CurrentUrl Query: {}", rawQuery); // 過濾沒有 query string // 還有過濾無法成對 keyValue 的 query, e.g. http://host/path?123 if (Objects.isNull(rawQuery) |...

Spring-boot Thymeleaf Html5 SAXParseException 解析錯誤

thymeleaf 解析 html5 出錯 <head> <meta charset="utf-8"> <meta http-equiv="X-UA-Compatible" content="IE=edge"> <meta name="viewport" content="width=device-width, initial-scale=1, shrink-to-fit=no"> <meta name="description" content=""> <meta name="author" content=""> <title>SB Admin - Start Bootstrap Template</title> <!-- Bootstrap core CSS--> <link href="../static/vendor/bootstrap/css/bootstrap.min.css" rel="stylesheet"> <!-- Custom fonts for this template--> <link href="../static/vendor/font-awesome/css/font-awesome.min.css" rel="stylesheet" type="text/css"> <!-- Page level plugin CSS--> <link href="../static/vendor/datatables/dataTables.bootstrap4.css" rel="stylesheet"> <!-- Custom styles for this template--> <li...

Google Compute Engine‎ - AccessDeniedExceptions 403

原因 打算從 instance 打包 logs 到 google cloud storage 發生了 AccessDeniedException: 403 Insufficient OAuth2 scope to perform this operation. , 看起來是 instance 沒有 storage 權限 解決 Reference: https://cloud.google.com/compute/docs/access/create-enable-service-accounts-for-instances#changeserviceaccountandscopes 重新設定 service account 權限 instance 上內建有 gcloud , 就直接用現有的工具查詢一下 instance 的 account. $ gsutil info 或者在本機直接 gcloud compute instances describe INSTANCE_NAMES Account: [alpha-number-compute@developer.gserviceaccount.com] Project: [our-project-name] 會看到 instance 的一些狀態, 接下來就簡單多了, 按照下列的說明, 要先 stop instance, 更改 storage scope 再重新 start 。 To change an instance's service account and access scopes, the instance must be temporarily stopped. To stop your instance, read the documentation for Stopping an instance. After changing the service account or access scopes, remember to restart the instance. # Stop Instance gcloud compute instances stop INSTANCE_NAMES # 設定 storage scope 為 full (Read, Write) gcloud co...