【问题标题】:How to implement connection pooling for PostgreSQL in Google Cloud Dataflow如何在 Google Cloud Dataflow 中为 PostgreSQL 实现连接池
【发布时间】:2021-03-17 02:40:14
【问题描述】:

#源配置

source_config = relational_db.SourceConfiguration(
             drivername=CONFIG['drivername'],
                     host=CONFIG['host'],
                     port=CONFIG['port'],
                     database=CONFIG['database'],
                     username=CONFIG['username'],
                     password=CONFIG['password']
                    )
# Target database table
customer_purchase_config = relational_db.TableConfiguration(
    name = 'customer_purchase',
    create_if_missing = False,
    primary_key_columns = ['purchaseId']
    )

使用 beam.Pipeline(options=options) 作为 p:

    res = (
        p
        | "Read data from PubSub"
        >> beam.io.ReadFromPubSub(subscription=SUB).with_output_types(bytes)
        |'Transformation' >> (beam.ParDo(PubSubToDict()))
    )
    customer_purchase = res | beam.ParDo(Customer_Purchase())
    customer_purchase | 'Writing to customer_purchase' >> relational_db.Write(
        source_config=source_config,
        table_config=customer_purchase_config
        )

因此,当我尝试使用这些配置时,我能够在 PostgreSQL 中插入和更新数据,但是当我收到大量输入峰值时,我的连接限制已达到,并且来自工作节点的重试次数正在增加,所以有什么办法可以定义连接池,以便我可以重用连接。

【问题讨论】:

  • 您在为您的 Postgre SQL 使用 Cloud SQL 吗?如果是,您需要记住这些连接limitations 存在于实例的内存中。此外,根据文档,您可以按照here 的描述设置连接池。对你有帮助吗?
  • 这个例子是用 Java 编写的,但可能会有所帮助。 stackoverflow.com/a/59402251/4756279

标签: python postgresql google-cloud-dataflow apache-beam


【解决方案1】:

查看这个pull request 示例,了解如何在 Dataflow 中管理连接池和批量写入 - Python SDK 到 Cloud SQL。

这个想法是利用DoFn.Setup() 方法,该方法在对象创建时运行一次,将为每个 python 解释器创建一个静态连接池。该池将在对象的生命周期内存在,然后您是否可以通过租用连接并将连接返回池以供重复使用来打开和关闭会话。

经过一些测试,问题不在于连接数,而是每秒创建和关闭的新连接数。每个新连接都需要启动 SQL 服务器上的一个线程来处理事务,这是一种反模式,会导致大量开销。假设您正在使用 Beam Nuggets,该库会创建一个新连接并在每个捆绑包之后处理它。如果你有小包,那么它会太快地创建新的连接。这在 SQL 服务器端造成了一个问题,因为对同一受影响记录的每个事务现在由不同的线程处理,这可能导致线程之间的行锁争用,并导致性能瓶颈,如Cloud SQL documentation 所述。

另外要记住的是键的数量和包的大小。决定捆绑包大小的因素之一是每个键可用的数据量。从 Pub/Sub 读取时,默认为 1024 个键,目的是帮助过滤重复消息。您可以使用 shuffle 步骤 (GBK) GroupByKey 重新键入一组密钥,例如 40 个。这也有助于增加包大小。如果实施正确,您应该将相同记录的相同事务组合到相同的键中,然后将其作为一个完整事务批处理到 Cloud SQL。这允许 Cloud SQL 优化性能并减少锁定行的线程开销。

我的建议是使用少于 1024 个键的 GroupByKey,这会导致更大的包大小并将受影响记录的每个事务组合在一起。然后,您可以使用patch 批处理到 CSQL,并具有处理连接池的额外好处。这会创建更少的新连接,从而减少 Cloud SQL 中的开销。

【讨论】:

    【解决方案2】:

    我在这里有一个用于生产的异步 python 池选项。效果非常好。

    Python Postgres psycopg2 ThreadedConnectionPool exhausted

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-08-05
      • 1970-01-01
      • 1970-01-01
      • 2018-05-03
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多