tmq.c 4.9 KB
Newer Older
L
Liu Jicong 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18
/*
 * Copyright (c) 2019 TAOS Data, Inc. <jhtao@taosdata.com>
 *
 * This program is free software: you can use, redistribute, and/or modify
 * it under the terms of the GNU Affero General Public License, version 3
 * or later ("AGPL"), as published by the Free Software Foundation.
 *
 * This program is distributed in the hope that it will be useful, but WITHOUT
 * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or
 * FITNESS FOR A PARTICULAR PURPOSE.
 *
 * You should have received a copy of the GNU Affero General Public License
 * along with this program. If not, see <http://www.gnu.org/licenses/>.
 */

#include <stdio.h>
#include <string.h>
#include <assert.h>
L
Liu Jicong 已提交
19
#include <time.h>
L
Liu Jicong 已提交
20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46
#include "taos.h"

static int running = 1;
static void msg_process(tmq_message_t* message) {
  tmqShowMsg(message);
}

int32_t init_env() {
  TAOS* pConn = taos_connect("localhost", "root", "taosdata", NULL, 0);
  if (pConn == NULL) {
    return -1;
  }

  TAOS_RES* pRes = taos_query(pConn, "create database if not exists abc1 vgroups 1");
  if (taos_errno(pRes) != 0) {
    printf("error in create db, reason:%s\n", taos_errstr(pRes));
    return -1;
  }
  taos_free_result(pRes);

  pRes = taos_query(pConn, "use abc1");
  if (taos_errno(pRes) != 0) {
    printf("error in use db, reason:%s\n", taos_errstr(pRes));
    return -1;
  }
  taos_free_result(pRes);

L
Liu Jicong 已提交
47 48 49 50 51
  pRes = taos_query(pConn, "create stable st1 (ts timestamp, k int) tags(a int)");
  /*if (taos_errno(pRes) != 0) {*/
    /*printf("failed to create super table 123_$^), reason:%s\n", taos_errstr(pRes));*/
    /*return -1;*/
  /*}*/
L
Liu Jicong 已提交
52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68
  taos_free_result(pRes);

  pRes = taos_query(pConn, "create table tu using st1 tags(1)");
  if (taos_errno(pRes) != 0) {
    printf("failed to create child table tu, reason:%s\n", taos_errstr(pRes));
    return -1;
  }
  taos_free_result(pRes);

  pRes = taos_query(pConn, "create table tu2 using st1 tags(2)");
  if (taos_errno(pRes) != 0) {
    printf("failed to create child table tu2, reason:%s\n", taos_errstr(pRes));
    return -1;
  }
  taos_free_result(pRes);


L
Liu Jicong 已提交
69
  const char* sql = "select * from st1";
L
Liu Jicong 已提交
70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115
  pRes = tmq_create_topic(pConn, "test_stb_topic_1", sql, strlen(sql));
  /*if (taos_errno(pRes) != 0) {*/
    /*printf("failed to create topic test_stb_topic_1, reason:%s\n", taos_errstr(pRes));*/
    /*return -1;*/
  /*}*/
  /*taos_free_result(pRes);*/
  taos_close(pConn);
  return 0;
}

tmq_t* build_consumer() {
  TAOS* pConn = taos_connect("localhost", "root", "taosdata", NULL, 0);
  assert(pConn != NULL);

  TAOS_RES* pRes = taos_query(pConn, "use abc1");
  if (taos_errno(pRes) != 0) {
    printf("error in use db, reason:%s\n", taos_errstr(pRes));
  }
  taos_free_result(pRes);

  tmq_conf_t* conf = tmq_conf_new();
  tmq_conf_set(conf, "group.id", "tg2");
  tmq_t* tmq = tmq_consumer_new(pConn, conf, NULL, 0);
  return tmq;

  tmq_list_t* topic_list = tmq_list_new();
  tmq_list_append(topic_list, "test_stb_topic_1");
  tmq_subscribe(tmq, topic_list);
  return NULL; 
}

tmq_list_t* build_topic_list() {
  tmq_list_t* topic_list = tmq_list_new();
  tmq_list_append(topic_list, "test_stb_topic_1");
  return topic_list;
}

void basic_consume_loop(tmq_t *tmq,
                        tmq_list_t *topics) {
  tmq_resp_err_t err;

  if ((err = tmq_subscribe(tmq, topics))) {
    fprintf(stderr, "%% Failed to start consuming topics: %s\n", tmq_err2str(err));
    printf("subscribe err\n");
    return;
  }
L
Liu Jicong 已提交
116
  int32_t cnt = 0;
L
Liu Jicong 已提交
117
  /*clock_t startTime = clock();*/
L
Liu Jicong 已提交
118
  while (running) {
L
Liu Jicong 已提交
119
    tmq_message_t *tmqmessage = tmq_consumer_poll(tmq, 500);
L
Liu Jicong 已提交
120 121
    if (tmqmessage) {
      cnt++;
L
Liu Jicong 已提交
122
      msg_process(tmqmessage);
L
Liu Jicong 已提交
123
      tmq_message_destroy(tmqmessage);
L
Liu Jicong 已提交
124 125
    /*} else {*/
      /*break;*/
L
Liu Jicong 已提交
126 127
    }
  }
L
Liu Jicong 已提交
128 129
  /*clock_t endTime = clock();*/
  /*printf("log cnt: %d %f s\n", cnt, (double)(endTime - startTime) / CLOCKS_PER_SEC);*/
L
Liu Jicong 已提交
130 131 132 133 134 135 136 137

  err = tmq_consumer_close(tmq);
  if (err)
    fprintf(stderr, "%% Failed to close consumer: %s\n", tmq_err2str(err));
  else
    fprintf(stderr, "%% Consumer closed\n");
}

L
Liu Jicong 已提交
138
void sync_consume_loop(tmq_t *tmq,
L
Liu Jicong 已提交
139 140 141 142 143 144
                  tmq_list_t *topics) {
  static const int MIN_COMMIT_COUNT = 1000;

  int msg_count = 0;
  tmq_resp_err_t err;

L
Liu Jicong 已提交
145
  if ((err = tmq_subscribe(tmq, topics))) {
L
Liu Jicong 已提交
146 147 148 149 150
    fprintf(stderr, "%% Failed to start consuming topics: %s\n", tmq_err2str(err));
    return;
  }

  while (running) {
L
Liu Jicong 已提交
151
    tmq_message_t *tmqmessage = tmq_consumer_poll(tmq, 500);
L
Liu Jicong 已提交
152 153 154 155 156
    if (tmqmessage) {
      msg_process(tmqmessage);
      tmq_message_destroy(tmqmessage);

      if ((++msg_count % MIN_COMMIT_COUNT) == 0)
L
Liu Jicong 已提交
157
        tmq_commit(tmq, NULL, 0);
L
Liu Jicong 已提交
158 159 160
    }
 }

L
Liu Jicong 已提交
161
  err = tmq_consumer_close(tmq);
L
Liu Jicong 已提交
162 163 164 165 166 167 168 169 170 171 172
  if (err)
    fprintf(stderr, "%% Failed to close consumer: %s\n", tmq_err2str(err));
  else
    fprintf(stderr, "%% Consumer closed\n");
}

int main() {
  int code;
  code = init_env();
  tmq_t* tmq = build_consumer();
  tmq_list_t* topic_list = build_topic_list();
L
Liu Jicong 已提交
173 174
  basic_consume_loop(tmq, topic_list);
  /*sync_consume_loop(tmq, topic_list);*/
L
Liu Jicong 已提交
175
}