【问题标题】:Resolving bottleneck on database connection in Dataflow pipeline解决数据流管道中数据库连接的瓶颈
【发布时间】:2022-11-18 15:39:33
【问题描述】:

我们有一个数据流流作业,它在 Pubsub 中使用消息,进行一些转换,并在 CloudSQL Postgres 实例上执行 DML(插入、更新、删除)。我们观察到瓶颈在数据库中。代码是用 Python 编写的,并使用 SQLAlchemy 作为与 Postgres 接口的库

我们观察到的常见问题是:

  1. 它最大化了允许的数据库连接,创建了多个连接池。
  2. 当有大量数据从 Pubsub 传入时,负责写入数据库的 DoFn 会抛出以下异常:
    Task was destroyed but it is pending! task: <Task pending name='Task-194770'...
    Task exception was never retrieved future: <Task finished name='Task-196602'...
    
    RuntimeError: aiohttp.client_exceptions.ClientResponseError: 429, message='Too Many Requests', url=URL('https://sqladmin.googleapis.com/sql/v1beta4/projects/.../instances/db-csql:generateEphemeralCert') [while running 'write_data-ptransform-48']
    

    Cloud SQL API 似乎在此处达到了速率限制。

    这些应该是我们理想的场景:

    1. 无论 Dataflow 创建的工作量和数量如何,我们在整个管道中应该只有一个 ConnectionPool(单例),具有静态连接数(最多 50 个分配给 Dataflow 作业,最大连接数为 200 个)在数据库中配置)。
    2. 在来自 Pubsub 的大量流量时刻,应该有某种机制来限制对数据库的传入请求的速率。或者不缩放负责写入数据库的 DoFn 的工作人员数量。

      你能推荐一种方法来完成这个吗?

      根据我的经验,单个全局连接池是不可能的,因为您无法将连接对象传递给工作人员(pickle/unpickle)。这是真的?

【问题讨论】:

  • 您是否在 DoFnsetup 方法中实例化了连接池?这是为每个工作人员创建连接池的推荐方法。然后必须在 DoFn 生命周期的 teardown 方法中关闭连接。
  • @MazlumTosun 是的,这就是我们所做的。然而,在大量流动数据的时刻,为了缓解背压,Dataflow 也在 write_to_db_dofn 中创建了很多 worker,以便它最大化数据库本身配置的允许连接。有没有办法在管道中静态设置特定步骤中允许的工作人员数量,比如 2,所以我们只能有可预测的最大连接数?
  • 由于您的问题集中在为您的两个要求找到set-up recommendations,因此将您的问题重定向到的更合适的论坛是Software Engineering StackExchange 论坛。

标签: python sqlalchemy google-cloud-dataflow apache-beam google-cloud-sql


【解决方案1】:

您应该尝试批处理对数据库的调用。伪代码看起来像这样(取自beam programming guide

class BufferDoFn(DoFn):
  BUFFER = BagStateSpec('buffer', EventCoder())
  IS_TIMER_SET = ReadModifyWriteStateSpec('is_timer_set', BooleanCoder())
  OUTPUT = TimerSpec('output', TimeDomain.REAL_TIME)

  def process(self,
              buffer=DoFn.StateParam(BUFFER),
              is_timer_set=DoFn.StateParam(IS_TIMER_SET),
              timer=DoFn.TimerParam(OUTPUT)):
    buffer.add(element)
    if not is_timer_set.read():
      timer.set(Timestamp.now() + Duration(seconds=10))
      is_timer_set.write(True)

  @on_timer(OUTPUT)
  def output_callback(self,
                      buffer=DoFn.StateParam(BUFFER),
                      is_timer_set=DoFn.StateParam(IS_TIMER_SET)):
    send_rpc(list(buffer.read()))
    buffer.clear()
    is_timer_set.clear()

原则上,您需要写一个splittable dofn并使用timers and states

【讨论】:

    猜你喜欢
    • 2021-12-29
    • 2020-11-18
    • 2013-01-18
    • 1970-01-01
    • 1970-01-01
    • 2019-07-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多