diff --git a/config/sys-dev.config b/config/sys-dev.config index e72c833..062f386 100644 --- a/config/sys-dev.config +++ b/config/sys-dev.config @@ -50,20 +50,6 @@ ]} %{pools, [ - % %% mysql连接池配置 - % {mysql_iot, - % [{size, 10}, {max_overflow, 20}, {worker_module, mysql}], - % [ - % {host, "47.111.101.3"}, - % {port, 3306}, - % {user, "root"}, - % {connect_mode, synchronous}, - % {keep_alive, true}, - % {password, "r3a-7Qrh#3Q"}, - % {database, "nannong_demo"} - % ] - % }, - % %% redis连接池 % {redis_pool, % [{size, 10}, {max_overflow, 20}, {worker_module, eredis}], diff --git a/rebar.config b/rebar.config index 73be5ab..4386c95 100644 --- a/rebar.config +++ b/rebar.config @@ -36,7 +36,6 @@ {ranch, ".*", {git, "https://github.com/ninenines/ranch.git", {tag, "2.2.0"}}}, {brod, ".*", {git, "https://github.com/kafka4beam/brod.git", {tag, "4.4.5"}}}, {jiffy, ".*", {git, "https://github.com/davisp/jiffy.git", {tag, "1.1.1"}}}, - {mysql, ".*", {git, "https://github.com/mysql-otp/mysql-otp", {tag, "1.8.0"}}}, {eredis, ".*", {git, "https://github.com/wooga/eredis.git", {tag, "v1.2.0"}}}, {emqtt, ".*", {git, "https://gitea.s5s8.com/anlicheng/emqtt.git", {tag, "v1.2"}}}, {gproc, ".*", {git, "https://github.com/uwiger/gproc.git", {tag, "0.9.1"}}}, diff --git a/rebar.lock b/rebar.lock index 7ef174e..fe8f776 100644 --- a/rebar.lock +++ b/rebar.lock @@ -38,10 +38,6 @@ {<<"kafka_protocol">>,{pkg,<<"kafka_protocol">>,<<"4.2.7">>},1}, {<<"metrics">>,{pkg,<<"metrics">>,<<"1.0.1">>},1}, {<<"mimerl">>,{pkg,<<"mimerl">>,<<"1.4.0">>},1}, - {<<"mysql">>, - {git,"https://github.com/mysql-otp/mysql-otp", - {ref,"caf5ff96c677a8fe0ce6f4082bc036c8fd27dd62"}}, - 0}, {<<"parse_trans">>, {git,"https://github.com/uwiger/parse_trans", {ref,"6f3645afb43c7c57d61b54ef59aecab288ce1013"}}, diff --git a/src/adapters/mysql/mysql_pool.erl b/src/adapters/mysql/mysql_pool.erl deleted file mode 100644 index cbceac5..0000000 --- a/src/adapters/mysql/mysql_pool.erl +++ /dev/null @@ -1,48 +0,0 @@ -%%%------------------------------------------------------------------- -%%% @author aresei -%%% @copyright (C) 2018, -%%% @doc -%%% -%%% @end -%%% Created : 29. 九月 2018 17:01 -%%%------------------------------------------------------------------- --module(mysql_pool). --author("aresei"). - -%% API --export([get_row/2, get_row/3, get_all/2, get_all/3]). --export([update/4, update_by/2, update_by/3, insert/4]). - -%% 从数据库中查找一行记录 --spec get_row(Pool :: atom(), Sql::binary()) -> {ok, Record::map()} | undefined. -get_row(Pool, Sql) when is_atom(Pool), is_binary(Sql) -> - poolboy:transaction(Pool, fun(ConnPid) -> mysql_provider:get_row(ConnPid, Sql) end). - --spec get_row(Pool :: atom(), Sql::binary(), Params::list()) -> {ok, Record::map()} | undefined. -get_row(Pool, Sql, Params) when is_atom(Pool), is_binary(Sql), is_list(Params) -> - poolboy:transaction(Pool, fun(ConnPid) -> mysql_provider:get_row(ConnPid, Sql, Params) end). - --spec get_all(Pool :: atom(), Sql::binary()) -> {ok, Rows::list()} | {error, Reason :: any()}. -get_all(Pool, Sql) when is_atom(Pool), is_binary(Sql) -> - poolboy:transaction(Pool, fun(ConnPid) -> mysql_provider:get_all(ConnPid, Sql) end). - --spec get_all(Pool :: atom(), Sql::binary(), Params::list()) -> {ok, Rows::list()} | {error, Reason::any()}. -get_all(Pool, Sql, Params) when is_atom(Pool), is_binary(Sql), is_list(Params) -> - poolboy:transaction(Pool, fun(ConnPid) -> mysql_provider:get_all(ConnPid, Sql, Params) end). - --spec insert(Pool :: atom(), Table :: binary(), Fields :: map() | list(), boolean()) -> - ok | {ok, InsertId :: integer()} | {error, Reason :: any()}. -insert(Pool, Table, Fields, FetchInsertId) when is_atom(Pool), is_binary(Table), is_list(Fields); is_map(Fields), is_boolean(FetchInsertId) -> - poolboy:transaction(Pool, fun(ConnPid) -> mysql_provider:insert(ConnPid, Table, Fields, FetchInsertId) end). - --spec update_by(Pool :: atom(), UpdateSql :: binary()) -> {ok, AffectedRows :: integer()} | {error, Reason :: any()}. -update_by(Pool, UpdateSql) when is_atom(Pool), is_binary(UpdateSql) -> - poolboy:transaction(Pool, fun(ConnPid) -> mysql_provider:update_by(ConnPid, UpdateSql) end). - --spec update_by(Pool :: atom(), UpdateSql :: binary(), Params :: list()) -> {ok, AffectedRows :: integer()} | {error, Reason :: any()}. -update_by(Pool, UpdateSql, Params) when is_atom(Pool), is_binary(UpdateSql) -> - poolboy:transaction(Pool, fun(ConnPid) -> mysql_provider:update_by(ConnPid, UpdateSql, Params) end). - --spec update(Pool :: atom(), Table :: binary(), Fields :: map(), WhereFields :: map()) -> {ok, AffectedRows::integer()} | {error, Reason::any()}. -update(Pool, Table, Fields, WhereFields) when is_atom(Pool), is_binary(Table), is_map(Fields), is_map(WhereFields) -> - poolboy:transaction(Pool, fun(ConnPid) -> mysql_provider:update(ConnPid, Table, Fields, WhereFields) end). \ No newline at end of file diff --git a/src/adapters/mysql/mysql_provider.erl b/src/adapters/mysql/mysql_provider.erl deleted file mode 100644 index e246124..0000000 --- a/src/adapters/mysql/mysql_provider.erl +++ /dev/null @@ -1,144 +0,0 @@ -%%%------------------------------------------------------------------- -%%% @author aresei -%%% @copyright (C) 2018, -%%% @doc -%%% -%%% @end -%%% Created : 29. 九月 2018 17:01 -%%%------------------------------------------------------------------- --module(mysql_provider). --author("aresei"). - -%% API --export([get_row/2, get_row/3, get_all/2, get_all/3]). --export([update/4, update_by/2, update_by/3, insert/4]). - -%% 从数据库中查找一行记录 --spec get_row(ConnPid :: pid(), Sql::binary()) -> {ok, Record::map()} | undefined. -get_row(ConnPid, Sql) when is_pid(ConnPid), is_binary(Sql) -> - logger:debug("[mysql_client] get_row sql is: ~p", [Sql]), - case mysql:query(ConnPid, Sql) of - {ok, Names, [Row | _]} -> - {ok, maps:from_list(lists:zip(Names, Row))}; - {ok, _, []} -> - undefined; - Error -> - logger:warning("[mysql_client] get error: ~p", [Error]), - undefined - end. - --spec get_row(ConnPid :: pid(), Sql::binary(), Params::list()) -> {ok, Record::map()} | undefined. -get_row(ConnPid, Sql, Params) when is_pid(ConnPid), is_binary(Sql), is_list(Params) -> - logger:debug("[mysql_client] get_row sql is: ~p, params: ~p", [Sql, Params]), - case mysql:query(ConnPid, Sql, Params) of - {ok, Names, [Row | _]} -> - {ok, maps:from_list(lists:zip(Names, Row))}; - {ok, _, []} -> - undefined; - Error -> - logger:warning("[mysql_client] get error: ~p", [Error]), - undefined - end. - --spec get_all(ConnPid :: pid(), Sql::binary()) -> {ok, Rows::list()} | {error, Reason :: any()}. -get_all(ConnPid, Sql) when is_pid(ConnPid), is_binary(Sql) -> - logger:debug("[mysql_client] get_all sql is: ~p", [Sql]), - case mysql:query(ConnPid, Sql) of - {ok, Names, Rows} -> - {ok, lists:map(fun(Row) -> maps:from_list(lists:zip(Names, Row)) end, Rows)}; - {error, Reason} -> - logger:warning("[mysql_client] get error: ~p", [Reason]), - {error, Reason} - end. - --spec get_all(ConnPid :: pid(), Sql::binary(), Params::list()) -> {ok, Rows::list()} | {error, Reason::any()}. -get_all(ConnPid, Sql, Params) when is_pid(ConnPid), is_binary(Sql), is_list(Params) -> - logger:debug("[mysql_client] get_all sql is: ~p, params: ~p", [Sql, Params]), - case mysql:query(ConnPid, Sql, Params) of - {ok, Names, Rows} -> - {ok, lists:map(fun(Row) -> maps:from_list(lists:zip(Names, Row)) end, Rows)}; - {error, Reason} -> - logger:warning("[mysql_client] get error: ~p", [Reason]), - {error, Reason} - end. - --spec insert(ConnPid :: pid(), Table :: binary(), Fields :: map() | list(), boolean()) -> - ok | {ok, InsertId :: integer()} | {error, Reason :: any()}. -insert(ConnPid, Table, Fields, FetchInsertId) when is_pid(ConnPid), is_binary(Table), is_map(Fields), is_boolean(FetchInsertId) -> - insert(ConnPid, Table, maps:to_list(Fields), FetchInsertId); -insert(ConnPid, Table, Fields, FetchInsertId) when is_pid(ConnPid), is_binary(Table), is_list(Fields), is_boolean(FetchInsertId) -> - {Keys, Values} = kvs(Fields), - - FieldSql = iolist_to_binary(lists:join(<<", ">>, Keys)), - Placeholders = lists:duplicate(length(Keys), <<"?">>), - ValuesPlaceholder = iolist_to_binary(lists:join(<<", ">>, Placeholders)), - - Sql = <<"INSERT INTO ", Table/binary, "(", FieldSql/binary, ") VALUES(", ValuesPlaceholder/binary, ")">>, - logger:debug("[mysql_client] insert sql is: ~p, params: ~p", [Sql, Values]), - case mysql:query(ConnPid, Sql, Values) of - ok -> - case FetchInsertId of - true -> - InsertId = mysql:insert_id(ConnPid), - {ok, InsertId}; - false -> - ok - end; - Error -> - Error - end. - --spec update_by(ConnPid :: pid(), UpdateSql :: binary()) -> {ok, AffectedRows :: integer()} | {error, Reason :: any()}. -update_by(ConnPid, UpdateSql) when is_pid(ConnPid), is_binary(UpdateSql) -> - logger:debug("[mysql_client] updateBySql sql: ~p", [UpdateSql]), - case mysql:query(ConnPid, UpdateSql) of - ok -> - AffectedRows = mysql:affected_rows(ConnPid), - {ok, AffectedRows}; - Error -> - Error - end. - --spec update_by(ConnPid :: pid(), UpdateSql :: binary(), Params :: list()) -> {ok, AffectedRows :: integer()} | {error, Reason :: any()}. -update_by(ConnPid, UpdateSql, Params) when is_pid(ConnPid), is_binary(UpdateSql) -> - logger:debug("[mysql_client] updateBySql sql: ~p, params: ~p", [UpdateSql, Params]), - case mysql:query(ConnPid, UpdateSql, Params) of - ok -> - AffectedRows = mysql:affected_rows(ConnPid), - {ok, AffectedRows}; - Error -> - Error - end. - --spec update(ConnPid :: pid(), Sql :: binary(), Fields :: map(), WhereFields :: map()) -> - {ok, AffectedRows::integer()} | {error, Reason::any()}. -update(ConnPid, Table, Fields, WhereFields) when is_pid(ConnPid), is_binary(Table), is_map(Fields), is_map(WhereFields) -> - %% 拼接set - {SetKeys, SetVals} = kvs(Fields), - SetKeys1 = lists:map(fun(K) when is_binary(K) -> <<"`", K/binary, "` = ?">> end, SetKeys), - SetSql = iolist_to_binary(lists:join(<<", ">>, SetKeys1)), - - %% 拼接where - {WhereKeys, WhereVals} = kvs(WhereFields), - WhereKeys1 = lists:map(fun(K) when is_binary(K) -> <<"`", K/binary, "` = ?">> end, WhereKeys), - WhereSql = iolist_to_binary(lists:join(<<" AND ">>, WhereKeys1)), - - Params = SetVals ++ WhereVals, - - Sql = <<"UPDATE ", Table/binary, " SET ", SetSql/binary, " WHERE ", WhereSql/binary>>, - logger:debug("[mysql_client] update sql is: ~p, params: ~p", [Sql, Params]), - case mysql:query(ConnPid, Sql, Params) of - ok -> - AffectedRows = mysql:affected_rows(ConnPid), - {ok, AffectedRows}; - Error -> - logger:error("[mysql_client] update sql: ~p, params: ~p, get a error: ~p", [Sql, Params, Error]), - Error - end. - --spec kvs(Fields :: map() | list()) -> {Keys :: list(), Values :: list()}. -kvs(Fields) when is_map(Fields) -> - kvs(maps:to_list(Fields)); -kvs(Fields) when is_list(Fields) -> - {Keys0, Values0} = lists:foldl(fun({K, V}, {Acc0, Acc1}) -> {[K|Acc0], [V|Acc1]} end, {[], []}, Fields), - {lists:reverse(Keys0), lists:reverse(Values0)}. \ No newline at end of file diff --git a/src/devtools/iot_mock.erl b/src/devtools/iot_mock.erl index 2848e52..cce8fa2 100644 --- a/src/devtools/iot_mock.erl +++ b/src/devtools/iot_mock.erl @@ -12,7 +12,6 @@ %% API -export([rsa_encode/1]). --export([insert_services/1]). -export([test_influxdb/0]). test_influxdb() -> @@ -29,20 +28,6 @@ test_influxdb() -> end) end, lists:seq(1, 100)). -insert_services(Num) -> - lists:foreach(fun(Id) -> - Res = mysql_pool:insert(mysql_iot, <<"micro_service">>, - #{ - <<"name">> => <<"微服务"/utf8, (integer_to_binary(Id))/binary>>, - <<"code">> => <<"1223423423423423"/utf8>>, - <<"type">> => 1, - <<"version">> => <<"v1.0">>, - <<"url">> => <<"https://www.baidu.com">>, - <<"detail">> => <<"这是一个关于测试的微服务"/utf8>> - }, false), - logger:debug("insert service result is: ~p", [Res]) - end, lists:seq(1, Num)). - rsa_encode(Data) when is_binary(Data) -> %% 读取相关配置 PublicPemFile = "/tmp/keys/public.pem", @@ -77,4 +62,4 @@ rsa_decode(EncData) when is_binary(EncData) -> PlainData = public_key:decrypt_private(EncData, PubKey), logger:debug("plain data is: ~p", [PlainData]), - ok. \ No newline at end of file + ok. diff --git a/src/host/iot_host.erl b/src/host/iot_host.erl index a923dfe..89ae9b6 100644 --- a/src/host/iot_host.erl +++ b/src/host/iot_host.erl @@ -185,7 +185,7 @@ init([UUID]) -> end, {ok, StateName, #state{host_id = HostId, uuid = UUID, has_session = false}}; undefined -> - logger:warning("[iot_host] host uuid: ~p, loaded from mysql failed", [UUID]), + logger:warning("[iot_host] host uuid: ~p, load host info failed", [UUID]), ignore end. diff --git a/src/iot.app.src b/src/iot.app.src index 2ff6c79..969489b 100644 --- a/src/iot.app.src +++ b/src/iot.app.src @@ -15,7 +15,6 @@ parse_trans, hackney, poolboy, - mysql, gproc, % gpb, mnesia,