【问题标题】:Writting data from pubsub to bigtable via cloud functions将数据从 pubsub 写入 bigtable 与云函数
【发布时间】:2020-03-26 11:37:35
【问题描述】:

我是云大表的初学者,在使用云函数将数据从 pub/sub 写入大表时遇到很大问题。

云函数从 pubsub 获取消息,但问题在于下一步,将其写入 bigtable。

消息在 python 脚本中创建并发送到 pub/sub。

消息的一个例子:

b'{"eda":2.015176,"temperature":33.39,"bvp":-0.49,"x_acc":-36.0,"y_acc":-38.0,"z_acc":-128.0,"heart_rate": 83.78,"iddevice":15.0,"timestamp":"2019-12-01T20:01:36.927Z"}'

为了将它写入 bigtable,我创建了一个表:

 from google.cloud import bigtable 
 from google.cloud.bigtable import column_family

 client = bigtable.Client(project="projectid", admin=True) 
 instance = client.instance("bigtableinstance")
 table = instance.table("bigtable1")
 print('Creating the {} table.'.format(table)) 
 print('Creating columnfamily cf1 with Max Version GC rule...')
 max_versions_rule = column_family.MaxVersionsGCRule(2)
 column_family_id = 'cf1'
 column_families = {column_family_id: max_versions_rule}
 if not table.exists():
     table.create(column_families=column_families)
     print("Table {} is created.".format(table)) 
 else:
     print("Table {} already exists.".format(table))

这没有问题。

现在我尝试使用 main 方法在云函数中使用以下 python 代码通过 pub/sub 将消息写入 bigtable:

import json
import base64
import os
from google.cloud import bigtable
from google.cloud.bigtable import column_family, row_filters


project_id = os.environ.get('projetid', 'UNKNOWN')
INSTANCE = 'bigtableinstance'
TABLE = 'bigtable1'

client = bigtable.Client(project=project_id, admin=True)
instance = client.instance(INSTANCE)

colFamily = "cf1"
def writeToBigTable(table, data):
#    Parameters row_key (bytes) – The key for the row being created.
#    Returns A row owned by this table.
        row_key = data[colFamily]['iddevice'].value.encode()
        row = table.row(row_key)
        for colFamily in data.keys():
            for key in data[colFamily].keys():
                row.set_cell(colFamily,
                                        key,
                                        data[colFamily][key])
        table.mutate_rows([row])
        return data

def selectTable():
    stage = os.environ.get('stage', 'dev')
    table_id = TABLE + stage
    table = instance.table(table_id)
    return table


def main(event, context):
    data = base64.b64decode(event['data']).decode('utf-8')
    print("DATA: {}".format(data))
    eda, temperature, bvp, x_acc, y_acc, z_acc, heart_rate, iddevice, timestamp = data.split(',')

    table = selectTable()

    data = {'eda': eda,
         'temperature': temperature,
         'bvp': bvp,
         'x_acc':x_acc,
         'y_acc':y_acc,
         'z_acc':z_acc,
         'heart_rate':heart_rate,
         'iddevice':iddevice,
         'timestamp':timestamp}
    writeToBigTable(table, data)
    print("Data Written: {}".format(data))

我尝试了不同的版本,但找不到解决方案。

感谢您的帮助。

一切顺利

多米尼克

【问题讨论】:

  • 因为你是BigTable的初学者,所以我要问你这个问题:你有高吞吐量的顾虑吗?
  • 不,没有高吞吐量。它是一个带有小数据集的原型。
  • 好的,请注意 BigTable 的成本 ;-)。每秒最多 1M 的写入,您可以为此使用 BigQuery 流写入。
  • 能否确保您使用的是最新版本的客户端库,然后使用 table. direct_row(row_key) 而不是 table.row(row_key)

标签: python google-cloud-functions google-cloud-bigtable


【解决方案1】:

我认为这行是错误的:

    row_key = data[colFamily]['iddevice'].value.encode()

您正在传递数据对象,但它没有“cf1”属性。您也不必对其进行编码。试试这个:

    row_key = data['iddevice']

你的 for 循环也会有同样的问题。我认为这才是你想要的

    for col in data.keys():
        row.set_cell(colFamily, key, data[key])

另外,我知道您只是在玩它,但是使用设备 ID 作为行键的唯一值最终会很糟糕。推荐的可能是将行键和日期或您的其他属性之一(取决于您的查询)并使用它作为您的行键。 Cloud Bigtable schema 上的文档很有帮助,codelab 使用了更真实的示例数据集,并介绍了如何为该示例选择模式。它在 Java 中,但您仍然可以导入数据并运行您自己的查询。

【讨论】:

    【解决方案2】:

    首先非常感谢您的帮助。

    我尝试使用您的代码推荐来修复它,但不幸的是,由于其他错误,它现在无法正常工作。

    AttributeError: 'DirectRow' 对象没有属性 'append'

    我猜这是在下面的代码行中

            row.set_cell(colFamily,
                         key,
                         data[key])
    

    我可以想象错误的来源是字符串“数据”的拆分

    eda, temperature, bvp, x_acc, y_acc, z_acc, heart_rate, iddevice, timestamp = data.split(',')
    

    例如eda 看起来像这样:

    "'eda':2.015176"
    

    这对我来说看起来很不对。

    特别是当我将它插入以下字典时:

     data = {'eda': eda,....}
    

    错误

    AttributeError: 'DirectRow' 对象没有属性 'append' 似乎是说,我想用 set_cell 处理的数据有问题。有说 set_cell 以行作为列表或直接行实例的任何其他迭代。不应该适合它吗?

    我尝试了使用列表的解决方法,但这似乎使情况变得更糟。

    client = bigtable.Client(project=project_id, admin=True)
    instance = client.instance(INSTANCE)
    
    colFamily = "cf1"
    def writeToBigTable(table, dat):
    
        row_key = "{}-{}".format(dat[16], dat[17])
        row = table.row(row_key)
        for n in range(len(dat)):
            row.set_cell(colFamily,
                         dat[n],
                         dat[n+9])
        table.mutate_rows([row])
        return dat
    
    def selectTable():
        stage = os.environ.get('stage', 'dev')
        table_id = TABLE + stage
        table = instance.table(table_id)
        return table
    
    
    def main(event, context):
        data = base64.b64decode(event['data']).decode('utf-8')
        print("DATA: {}".format(data))
        var_1, eda, var_2, temperature, var_3, bvp, var_4, x_acc, var_5, y_acc, var_6, z_acc, var_7, heart_rate, var_8, iddevice, var_9, timestamp = data.replace(':',',').split(',')
    
        table = selectTable(); dat = [var_1, var_2, var_3, var_4, var_5, var_6, var_7, var_8, var_9, eda, temperature, bvp, x_acc, y_acc, z_acc, heart_rate, iddevice, timestamp]; 
    
    #   data = {'eda': eda,
    #         'temperature': temperature,
    #         'bvp': bvp,
    #         'x_acc':x_acc,
    #         'y_acc':y_acc,
    #         'z_acc':z_acc,
    #         'heart_rate':heart_rate,
    #         'iddevice':iddevice,
    #         'timestamp':timestamp}
        writeToBigTable(table, dat)
        print("Data Written: {}".format(data))
    

    我真的很难解决这个问题,并且不知道如何解决它。

    【讨论】:

    • 你以前的方式很好,你不需要做一个列表。不知道为什么你会得到那个附加错误,并且可以进一步研究它。此外,在此站点上,我建议您回复答案并编辑您的问题,而不是使用更新的代码创建新答案。这让其他人更容易跟进和提供帮助。
    猜你喜欢
    • 2022-12-12
    • 1970-01-01
    • 2020-03-04
    • 2018-12-15
    • 1970-01-01
    • 2021-09-07
    • 2020-11-26
    • 1970-01-01
    • 2017-09-02
    相关资源
    最近更新 更多