【问题标题】:Passing MySQL fetched rows to thread pool in C将 MySQL 获取的行传递给 C 中的线程池
【发布时间】:2022-09-25 08:14:38
【问题描述】:

我想同时处理从 MySQL 数据库中获取的数据。我将数据传递给每个线程进程(无需考虑线程安全;行在每个线程中独立处理):

#include <mysql.h>
#include <stdio.h>
#include <stdlib.h>
#include <stdint.h>
#include <pthread.h>
#include \"thpool.h\" // https://github.com/Pithikos/C-Thread-Pool

#define THREADS 10

struct fparam
{
  int id;
  char *data;
};

void process(void *arg)
{
  struct fparam *args = arg;
  // Processing ID and Data here
  printf(\"%d - %s\\n\", args->id, args->data);
}

int main(int argc, char **argv)
{
  threadpool thpool = thpool_init(THREADS);

  // MySQL connection

  MYSQL_RES *result = mysql_store_result(con);

  int num_fields = mysql_num_fields(result);
  struct fparam items[100]; // 100 is for the representation

  MYSQL_ROW row;
  int i = 0;
  while ((row = mysql_fetch_row(result)))
  {
    items[i].id = atoi(row[0]);
    items[i].data = row[1];
    thpool_add_work(thpool, process, (void *)(&items[i]));
    i++;
  }

  mysql_free_result(result);
  mysql_close(con);

  thpool_wait(thpool);
  thpool_destroy(thpool);

  exit(0);
}

当有很多行时,items 变得太大而无法放入内存(不仅仅是堆)。

如何限制存储在内存中的行数并在处理后删除它们?

我认为一个关键问题是我们不知道process 函数是否更快或从数据库中获取行。

  • “不仅仅是堆”是什么意思?你是说你不想使用堆?如果是这样,为什么不呢?
  • @kaylum 抱歉,我后来添加了它以避免在代码中不使用 malloc 造成混淆。我对堆或堆栈都很好。
  • 你是说数据行太多,连动态内存都太大了?在这种情况下,您需要在主线程和池线程之间进行同步,以便在池线程准备好接收它们时协调仅读取更多行。例如,使用计数信号量。
  • 听起来您需要在结果集(可能是巨大的 #/rows)和线程池(有限的 #/worker 线程)之间实现一个队列。
  • 如您所知,任何时候系统可能收到的数据多于它无法及时提供的服务,您应该考虑使用某种“队列”。以下是几个示例(您可以通过简单的 Google 搜索找到更多示例):log2base2.com/data-structures/queue/queue-data-structure.htmlprogramiz.com/dsa/circular-queue 等等。您的工作线程读取下一个可用项目 (\"dequeue\") 并为其提供服务。即使“服务”可以并行发生,您的“出队”也可能需要一个锁。

标签: c


【解决方案1】:

使用queue,一个列表,您可以在其中添加项目并将它们从另一端取出。

你可以自己写; linked list 可以用作队列,将项目添加到一端并从另一端移除它们。或者使用现有的实现,例如the one provided by GLib

【讨论】:

  • 在 cmets 中建议使用queue,但它并不能简单地解决我的问题。我不能将 MySQL 行直接传递到线程中。如何存储线程要读取的 n 行?
【解决方案2】:

您不需要在您的场景中创建新队列,因为thpool_init(THREADS) 已经为您提供了一个,而thpool_add_work 提供了这个内部队列并且它在获取时增长。如果从数据库中获取速度很快但处理速度很慢,则必须将获取新行的速度限制在合理的范围内,以便它们适合内存。查看"thpool.h" 的文档,有这个函数thpool_num_threads_working(threadpool)。它将返回工作线程的数量,以便您以最简单的形式定义THREADS,只要有至少一个空闲线程可用(类似于while(thpool_num_threads_working(thpool) &lt; THREADS)),您就希望获取新行。

考虑到性能原因,您应该考虑预取一些行以在任何线程完成其工作以便能够在不等待的情况下提供数据时让您的数据已经存在。这些行中有多少可以等待处理它取决于可用的内存。当获取可能非常耗时但处理速度很快时,考虑这一点更为重要。

你使用items[i] 的方式也会在这里产生一个问题,因为i 可以无限溢出items[100] 数组。如果我们放在那里非常强壮假设线程以与启动相同的顺序完成它们的工作,您可以重置i,这样items 特定索引将被重新用于新行(作为一种循环缓冲区)。不幸的是,现在恐怕C-Thread-Pool 不支持识别哪个特定线程已完成其工作(以及不再需要哪些相应数据)。如果您需要在这里 100% 安全,我会考虑两种可能的解决方案。扩展 C-Thread-Pool 与线程状态与它的作业 ID 验证或分批处理行,这样您就可以提供所有线程(每个线程一个作业到队列),然后等到所有线程完成他们的工作,然后再次将它们一起提供下一个要处理的一批行。

并且记得在使用thpool_add_workthpool_init 时检查错误。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2011-06-01
    • 1970-01-01
    • 1970-01-01
    • 2019-07-19
    • 1970-01-01
    • 2016-11-09
    • 1970-01-01
    • 2015-03-21
    相关资源
    最近更新 更多