騰訊云 CKafka 作為大資料架構中的關鍵組件,起到了資料聚合,流量削峰,訊息管道的作用,在 CKafka 上下游中的資料流轉中有各種優秀的開源解決方案,如 Logstash,File Beats,Spark,Flink 等等,本文將帶來一種新的解決方案:Serverless Function,其在學習成本,維護成本,擴縮容能力等方面相對已有開源方案將有優異的表現,
作者簡介:許文強,騰訊云 Ckafka 核心研發,精通 Kafka 及其周邊生態,對 Serverless,訊息佇列等領域有較深的理解,專注于 Kafka 在公有云多租戶和大規模集群場景下的性能分析和優化、及云上訊息佇列 serverless 化的相關探索,
Tencent Cloud Kafka 介紹
Tencent Cloud Kafka 是基于開源 Kafka 引擎研發的適合大規模公有云部署的 Cloud Kafka,是一款適合公有云部署、運行、運維的分布式的、高可靠、高吞吐和高可擴展的訊息佇列系統,它 100% 兼容開源的 Kafka API,目前主要支持開源的 0.9, 0.10, 1.1.1, 2.4.2 四個大版本 ,并提供向下兼容的能力,
目前 Tencent Cloud Kafka 維護了超過 4000+ 節點的集群,每日吞吐的訊息量超過 9 萬億+條,峰值帶寬達到了 800GB+/s, 堆積資料達到了 20PB+,是一款集成了租戶隔離、限流、鑒權、安全、資料監控告警、故障快速切換、跨可用區容災等等一系列特性的,歷經大流量檢驗的、可靠的公有云上 Kafka 集群,
什么是資料流轉
CKafka 作為一款高吞吐,高可靠的訊息佇列引擎,需要承接大量資料的流入和流出,資料流動的這一程序我們稱之它為資料流轉,而在處理資料的流入和流出程序中,會有很多成熟豐富的開源的解決方案,如 Logstash,Spark,Fllink等,從簡單的資料轉儲,到復雜的資料清洗,過濾,聚合等,都有現成的解決方案,
如圖所示,在 Kafka 上下游生態圖中,CKafka 處于中間層,起到資料聚合,流量削峰,訊息管道的作用,圖左和圖上是資料寫入的組件概覽,圖右和圖下是下游流式資料處理方案和持久化存盤引擎,這些構成了 Kafka 周邊的資料流動的生態,

資料流轉新方案: Serverless Function
下圖是流式計算典型資料流動示意圖,其中承接資料流轉方案的是各種開源解決方案,單純從功能和性能的角度來講,開源解決方案都有很優秀的表現,

而從學習成本,維護成本,金錢成本,擴縮容能力等角度來看,這些開源方案還是有欠缺的,怎么說呢?開源方案的缺點主要在于如下三點:
- 學習成本
- 調優、維護、解決問題的成本
- 擴縮容能力
以 Logstash 為例,它的入門使用學習門檻不高,進階使用有一定的成本,主要包括眾多 release 版本的使用成本,引數調優和故障處理成本,后續的維護成本(行程可用性,單機的負載處理)等,如果用流式計算引擎,如 spark 和 flink,其雖然具有分布式調度能力和即時的資料處理能力,但是其學習門檻和后期的集群維護成本,將大大提高,
來看 Serverless Function 是怎么處理資料流轉的,如圖所示,Serverless Function 運行在資料的流入和流出的處理層的位置,代替了開源的解決方案,Serverless Function 是以自定義代碼的形式來實作資料清洗、過濾、聚合、轉儲等能力的,它具有學習成本低、無維護成本、自動擴縮容和按量計費等優秀特性,

接下來我們來看一下 Serverless Function 是怎么實作資料流轉的,并且了解一下其底層的運行機制及其優勢,
Serverless Function 實作資料流轉
首先來看一下怎么使用 Serverless Function 實作 Kafka To Elasticsearch 的資料流轉,下面以 Function 事件觸發的方式來說明 Function 是怎么實作低成本的資料清洗、過濾、格式化和轉儲的:
在業務錯誤日志采集分析的場景中,會將機器上的日志資訊采集并發送到服務端,服務端選擇 Kafka 作為訊息中間件,起到資料可靠存盤,流量削峰的作用,為了保存長時間的資料(月,年),一般會將資料清洗、格式化、過濾、聚合后,存盤到后端的分布式存盤系統,如 HDFS,HBASE,Elasticsearch 中,
以下代碼段分為三部分:資料源的訊息格式,處理后的目標訊息格式,功能實作的 Function 代碼段
- 源資料格式:
{
"version": 1,
"componentName": "trade",
"timestamp": 1595944295,
"eventId": 9128499,
"returnValue": -1,
"returnCode": 101103,
"returnMessage": "return has no deal return error[錯誤:缺少**c引數][seqId:u3Becr8iz*]",
"data": [],
"seqId": "@kibana-highlighted-field@u3Becr8iz@/kibana-highlighted-field@*"
}
- 目標資料格式:
{
"timestamp": "2020-07-28 21:51:35",
"returnCode": 101103,
"returnError": "return has no deal return error",
"returnMessage": "錯誤:缺少**c引數",
"requestId": "u3Becr8iz*"
}
- Function 代碼
Function 實作的功能是將資料從源格式,通過清洗,過濾,格式化轉化為目標資料格式,并轉儲到 Elasticsearch,代碼的邏輯很簡單:CKafka 收到訊息后,觸發了函式的執行,函式接收到資訊后會執行 convertAndFilter 函式的過濾,重組,格式化操作,將源資料轉化為目標格式,最后資料會被存盤到 Elasticsearch,
#!/usr/bin/python
# -*- coding: UTF-8 -*-
from datetime import datetime
from elasticsearch import Elasticsearch
from elasticsearch import helpers
esServer = "http://172.16.16.53:9200" # 修改為 es server 地址+埠 E.g. http://172.16.16.53:9200
esUsr = "elastic" # 修改為 es 用戶名 E.g. elastic
esPw = "PW123" # 修改為 es 密碼 E.g. PW2312321321
esIndex = "pre1" # es 的 index 設定
# ... or specify common parameters as kwargs
es = Elasticsearch([esServer],
http_auth=(esUsr, esPw),
sniff_on_start=False,
sniff_on_connection_fail=False,
sniffer_timeout=None)
def convertAndFilter(sourceStr):
target = {}
source = json.loads(sourceStr)
# 過濾掉returnCode=0的日志
if source["returnCode"] == 0:
return
dateArray = datetime.datetime.fromtimestamp(source["timestamp"])
target["timestamp"] = dateArray.strftime("%Y-%m-%d %H:%M:%S")
target["returnCode"] = source["returnCode"]
message = source["returnMessage"]
message = message.split("][")
errorInfo = message[0].split("[")
target["returnError"] = errorInfo[0]
target["returnMessage"] = errorInfo[1]
target["requestId"] = message[1].replace("]", "").replace("seqId:", "")
return target
def main_handler(event, context):
# 獲取 event Records 欄位并做轉化操作 資料結構 https://cloud.tencent.com/document/product/583/17530
for record in event["Records"]:
target = convertAndFilter(record)
action = {
"_index": esIndex,
"_source": {
"msgBody": target # 獲取 Ckafka 觸發器 msgBody
}
}
helpers.bulk(es, action)
return ("successful!")
看到這里,大家可能會發現,這個代碼段平時是處理單機的少量資料的腳本是一樣的,就是做轉化,轉儲,很簡單,其實很多分布式的系統做的系統從微觀的角度看,其實就是做的這么簡單的事情,分布式框架本身做的更多的是分布式調度,分布式運行,可靠性,可用性等等作業,細化到執行單元,功能其實和上面的代碼段是一樣的,
從宏觀來看,Serverless Function 做的事情和分布式計算框架 Spark, Flink 等做的事情是一樣的,都是調度,執行基本的執行單元,處理業務邏輯,區別在于用開源的方案,需要使用方去學習,使用,維護運行引擎,而 Serverless Function 則是平臺來幫用戶做這些事情,
接下來我們來看 Serverless Function 在底層是怎么去支持這些功能的,來看一下其底層的運行機制,如圖所示:

Function 作為一個代碼片段,提交給平臺以后,需要有一種觸發函式運行的方式,目前主要有如下三種:事件觸發、定時觸發和主動觸發,
在上面的例子中,我們是以事件觸發為例的,當訊息提交到 Kafka,就會觸發函式的運行,此時 Serverless 調度運行平臺就會調度底層的 Container 并發去執行函式,并執行函式的邏輯,此時關于 Container 的并發度是由系統自動調度,自動計算的,當 Kafka 的源資料多的時候,并發量就大,當資料少的時候,相應的就會較少并發數,因為函式是以運行時長計費的,當源訊息資料量少的時候,并發量小,自然運行時長就少,自然所需付出的資金成本就降下來,
在函式執行程序當中,函式的可靠性運行,自動擴縮容調度,并發度等都是用戶不需要關心的,用戶需要 Cover 的只是函式代碼段的可運行,無 BUG,這對于研發人員的精力投入成本就降低很多,
值得一談的是,在開發語言方面,開源方案只支持其相對應的語言,如 Logstash 的嵌入腳本用的是 ruby,spark 主要支持java,scala,python 等,而 Serverless Function 支持的是幾乎業界常見到的開發語言,包括不限于 java,golang,python,node JS,php 等等,這點就可以讓研發人員用其熟悉的語言去解決資料流轉問題,這在無形中就減少了很多代碼出錯和出問題的機會,
Serverless Function 在資料流轉場景的優勢
下面我們來統一看一下 Serverless Function 和開源的方案的主要區別及優勢,如圖5所示,和開源方案相比,在非實時的資料流轉場景中,Serverless Function 相對現有的開源方案,它具有的優勢幾乎是壓倒性的,從功能和性能的角度,它在批式計算(實時)的場景中是完全可以滿足的,但是它相對開源方案在學習成本,運維成本幾乎可以忽略,其動態擴縮容,按需付費,毫秒級付費對于資金成本的投入也是非常友好的,

用一句話總結就是:Serverless Function 能用一段熟悉的語言撰寫一小段代碼去銜接契合流式計算中的資料流轉,
Serverless Function 在批式計算場景的展望
隨著流式計算的發展,慢慢演化出了批量計算 (batch computing)、流式計算 (stream computing)、互動計算 (interactive computing)、圖計算 (graph computing) 等方向,而架構師在業務中選擇批式計算或者流式計算,其核心是希望按需使用批式計算或流式計算,以取得在延時、吞吐、容錯、成本投入等方面的平衡,在使用者看來,批式處理可以提供精確的批式資料視圖,流式處理可以提供近實時的資料視圖,而在批式處理當中,或者說在未來的批式處理和流式處理的底層技術的合流程序中,Lambda 架構是其發展的必然路徑,
Serverless Function 以其按需使用,自動擴縮容及近乎無限的橫向擴容能力給現階段的批式處理提供了一種選擇,并且在未來批流一體化的程序中,未來可期,

One More Thing
立即體驗騰訊云 Serverless Demo,領取 Serverless 新用戶禮包 ?? serverless/start
歡迎訪問:Serverless 中文網!
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/699.html
標籤:其他
