【问题标题】:Fastest way to insert many rows of data?插入多行数据的最快方法?
【发布时间】:2021-10-30 07:50:54
【问题描述】:

每 4 秒,我必须存储 32,000 行数据。这些行中的每一行都包含一个时间戳值和 464 个双精度值。时间戳的列名是time,精度值的列名依次增加为channel1channel2、...和channel 464

我建立连接如下:

CONNECTION = f"postgres://{username}:{password}@{host}:{port}/{dbname}"#?sslmode=require"
self.TimescaleDB_Client = psycopg2.connect(CONNECTION)

然后,我使用以下内容验证 TimescaleDB 扩展:

def verifyTimeScaleInstall(self): 
    try: 
        sql_query = "CREATE EXTENSION IF NOT EXISTS timescaledb CASCADE;"
        cur = self.TimescaleDB_Client.cursor()
        cur.execute(sql_query)
        cur.close()
        self.TimescaleDB_Client.commit()
    except: 
        self.timescaleLogger.error("An error occurred in verifyTimeScaleInstall")
        tb = traceback.format_exc()
        self.timescaleLogger.exception(tb) 
        return False

然后,我使用以下内容为我的数据创建一个 hyptertable:

def createRAWDataTable(self): 
    try: 
        cur = self.TimescaleDB_Client.cursor()
        self.query_create_raw_data_table = None
        for channel in range(self.num_channel) :
            channel = channel + 1 
            if self.query_create_raw_data_table is None: 
                self.query_create_raw_data_table = f"CREATE TABLE IF NOT EXISTS raw_data (time TIMESTAMPTZ NOT NULL, channel{channel} REAL"
            else: 
                self.query_create_raw_data_table = self.query_create_raw_data_table + f", channel{channel} REAL"
        self.query_create_raw_data_table = self.query_create_raw_data_table + ");"
        self.query_create_raw_data_hypertable = "SELECT create_hypertable('raw_data', 'time');"
        cur.execute(self.query_create_raw_data_table)
        cur.execute(self.query_create_raw_data_hypertable)
        self.TimescaleDB_Client.commit()
        cur.close()
    except:         
        self.timescaleLogger.error("An error occurred in createRAWDataTable")
        tb = traceback.format_exc()
        self.timescaleLogger.exception(tb) 
        return False

然后我使用以下命令将数据插入到超表中:

def insertRAWData(self, seconds): 
    try: 
        insert_start_time = datetime.now(pytz.timezone("MST"))
        current_time = insert_start_time
        num_iterations = seconds * self.fs
        time_increment = timedelta(seconds=1/self.fs)

        raw_data_query = self.query_insert_raw_data

        dtype = "float32"
        matrix = np.random.rand(self.fs*seconds,self.num_channel).astype(dtype)
        cur = self.TimescaleDB_Client.cursor()
        data = list()
        for iteration in range(num_iterations): 
            raw_data_row = matrix[iteration,:].tolist() #Select a particular row and all columns
            time_string = current_time.strftime("%Y-%m-%d %H:%M:%S.%f %Z")
            raw_data_values = (time_string,)+tuple(raw_data_row)
            data.append(raw_data_values)
            current_time = current_time + time_increment

        start_time = time.perf_counter()
        psycopg2.extras.execute_values(
            cur, raw_data_query, data, template=None, page_size=100
        )
        print(time.perf_counter() - start_time)
        self.TimescaleDB_Client.commit()
        cur.close()

    except: 
        self.timescaleLogger.error("An error occurred in insertRAWData")
        tb = traceback.format_exc()
        self.timescaleLogger.exception(tb) 
        return False

我在上述代码中引用的 SQL 查询字符串是从以下获取的:

def getRAWData_Query(self): 
    try: 
        self.query_insert_raw_data = None
        for channel in range(self.num_channel): 
            channel = channel + 1
            if self.query_insert_raw_data is None: 
                self.query_insert_raw_data = f"INSERT INTO raw_data (time, channel{channel}"
            else: 
                self.query_insert_raw_data = self.query_insert_raw_data + f", channel{channel}"
        self.query_insert_raw_data = self.query_insert_raw_data + ") VALUES %s;"
        return self.query_insert_raw_data
    except: 
        self.timescaleLogger.error("An error occurred in insertRAWData_Query")
        tb = traceback.format_exc()
        self.timescaleLogger.exception(tb) 
        return False

如您所见,我使用psycopg2.extras.execute_values() 插入值。据我了解,这是插入数据的最快方法之一。但是,我插入这些数据大约需要 80 秒。它与12 cores/24 threadsSSDs256GB of RAM 在一个相当不错的系统上。这可以更快地完成吗?只是看起来很慢。

我想使用TimescaleDB 并正在评估它的性能。但我希望在 2 秒左右的时间内写出来,这样才能被接受。

编辑 我曾尝试使用pandas 来执行插入,但耗时更长,大约为 117 秒。以下是我使用的函数。

def insertRAWData_Pandas(self, seconds): 
    try: 
        insert_start_time = datetime.now(pytz.timezone("MST"))
        current_time = insert_start_time
        num_iterations = seconds * self.fs
        time_increment = timedelta(seconds=1/self.fs)

        raw_data_query = self.query_insert_raw_data
        dtype = "float32"
        matrix = np.random.rand(self.fs*seconds,self.num_channel).astype(dtype)
        pd_df_dict = {}
        pd_df_dict["time"] = list()
        for iteration in range(num_iterations):
            time_string = current_time.strftime("%Y-%m-%d %H:%M:%S.%f %Z") 
            pd_df_dict["time"].append(time_string)
            current_time = current_time + time_increment
        for channel in range(self.num_channel): 
            pd_df_dict[f"channel{channel}"] = matrix[:,channel].tolist()



        start_time = time.perf_counter()
        pd_df = pd.DataFrame(pd_df_dict)
        pd_df.to_sql('raw_data', self.engine, if_exists='append')
        print(time.perf_counter() - start_time)

    except: 
        self.timescaleLogger.error("An error occurred in insertRAWData_Pandas")
        tb = traceback.format_exc()
        self.timescaleLogger.exception(tb) 
        return False

edit 我尝试使用CopyManager,它似乎在大约 74 秒时产生了最好的结果。然而,仍然不是我想要的。

def insertRAWData_PGCOPY(self, seconds): 
    try: 
        insert_start_time = datetime.now(pytz.timezone("MST"))
        current_time = insert_start_time
        num_iterations = seconds * self.fs
        time_increment = timedelta(seconds=1/self.fs)
        dtype = "float32"
        matrix = np.random.rand(num_iterations,self.num_channel).astype(dtype)
        data = list()
        for iteration in range(num_iterations): 
            raw_data_row = matrix[iteration,:].tolist() #Select a particular row and all columns
            #time_string = current_time.strftime("%Y-%m-%d %H:%M:%S.%f %Z")
            raw_data_values = (current_time,)+tuple(raw_data_row)
            data.append(raw_data_values)
            current_time = current_time + time_increment
        channelList = list()
        for channel in range(self.num_channel): 
            channel = channel + 1
            channelString = f"channel{channel}"
            channelList.append(channelString)
        channelList.insert(0,"time")
        cols = tuple(channelList)

        start_time = time.perf_counter()
        mgr = CopyManager(self.TimescaleDB_Client, 'raw_data', cols)
        mgr.copy(data)
        self.TimescaleDB_Client.commit()
        print(time.perf_counter() - start_time)

    except: 
        self.timescaleLogger.error("An error occurred in insertRAWData_PGCOPY")
        tb = traceback.format_exc()
        self.timescaleLogger.exception(tb) 
        return False

我尝试修改postgresql.conf 中的以下值。没有明显的性能提升。

wal_level = minimal
fsync = off
synchronous_commit = off
wal_writer_delay = 2000ms
commit_delay = 100000

我已尝试根据以下 cmets 之一在我的 createRawDataTable() 函数中使用以下内容修改块大小。但是,插入时间没有改善。考虑到我没有积累数据,也许这也是可以预料的。数据库中的数据只有几个样本,在我的测试过程中可能最多 1 分钟。

self.query_create_raw_data_hypertable = "SELECT create_hypertable('raw_data', 'time', chunk_time_interval => INTERVAL '3 day',if_not_exists => TRUE);"

编辑 对于阅读本文的任何人,我能够在大约 0.5 秒内为 MongoDB 腌制并插入一个 32000x464 float32 numpy 矩阵,这是我的最终解决方案。在这种情况下,也许 MongoDB 在这种工作负载上做得更好。

【问题讨论】:

  • 如果您尝试不同的批量大小会很棒,因为您将根据需要将结果刷新到数据库中或多或少地分配。我还鼓励您尝试进行更多并行化,并将比率与更少的大批量或更多的小批量进行比较。找到正确的平衡点可能很棘手,但是通过一些参数场景,您将朝着正确的方向前进。即使在 Raspberry PI 上对 Timescale 进行基准测试,我每秒也能获得 20k+ 行。 ideia.me/time-series-benchmark-timescaledb-raspberry-pi

标签: python python-3.x timescaledb


【解决方案1】:

最快的方法之一是

  • 首先为要插入数据库的数据创建一个 pandas 数据框
  • 然后使用数据框将您的数据批量插入数据库中

您可以通过以下方式做到这一点:How to write data frame to postgres?

【讨论】:

  • 谢谢,不过时间好像增加到了 117 秒。
【解决方案2】:

我有两个初步建议可能有助于提高整体性能。

  1. 您创建的默认超表会将您的数据“分块”为 7 天(这意味着在给定参数的情况下,每个分块将包含大约 4,838,400,000 行数据)。由于您的数据是如此精细,您可能希望使用不同的块大小。查看docs here 以获取有关可选chunk_time_interval 参数的信息。更改卡盘大小应该有助于提高插入和查询速度,如果以后需要,它还可以为您提供更好的压缩性能。

  2. 正如上述个人所说,使用批量插入也应该有所帮助。如果您还没有查看此stock data tutorial,我强烈推荐它。使用pgcopy 及其函数CopyManager 可以帮助更快地插入df 对象。

希望这些信息对您的情况有所帮助!

披露:我是 Timescale 团队的一员?

【讨论】:

  • 你建议增加还是减少chunk_time_interval?
  • 我会减少间隔,以便在每个块中分区的行更少。希望这对您使用 CopyManager 功能获得的 4 秒有所帮助!
猜你喜欢
  • 1970-01-01
  • 2016-02-06
  • 2012-11-27
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-08-12
  • 1970-01-01
相关资源
最近更新 更多