【发布时间】: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