Skip to content

Commit 0ac8f09

Browse files
committed
Merge branch 'development'
* development: [Release] version 1.8.4 Update LICENSE.py Update setup.py Fix logging settings [Bugfix] Unable to write null values into DB
2 parents 07541c7 + d04a6bd commit 0ac8f09

9 files changed

Lines changed: 46 additions & 38 deletions

File tree

LICENSE.txt

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
MIT License
22

3-
Copyright (c) [2019] [esaki01]
3+
Copyright (c) [2021] [esakik]
44

55
Permission is hereby granted, free of charge, to any person obtaining a copy
66
of this software and associated documentation files (the "Software"), to deal

beam_mysql/__init__.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,3 @@
11
"""Beam - MySQL Connector information and utilities."""
22

3-
__version__ = "1.8.3"
3+
__version__ = "1.8.4"

beam_mysql/connector/client.py

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
"""A client of mysql."""
22

3-
import logging
3+
from logging import INFO, getLogger
44
from typing import Dict
55
from typing import Generator
66
from typing import List
@@ -13,6 +13,9 @@
1313
_SELECT_STATEMENT = "SELECT"
1414
_INSERT_STATEMENT = "INSERT"
1515

16+
logger = getLogger(__name__)
17+
logger.setLevel(INFO)
18+
1619

1720
class MySQLClient:
1821
"""A mysql client object."""
@@ -43,7 +46,7 @@ def record_generator(self, query: str, dictionary=True) -> Generator[Dict, None,
4346

4447
try:
4548
cur.execute(query)
46-
logging.info(f"Successfully execute query: {query}")
49+
logger.info(f"Successfully execute query: {query}")
4750

4851
for record in cur:
4952
yield record
@@ -74,7 +77,7 @@ def counts_estimator(self, query: str) -> int:
7477

7578
try:
7679
cur.execute(count_query)
77-
logging.info(f"Successfully execute query: {count_query}")
80+
logger.info(f"Successfully execute query: {count_query}")
7881

7982
record = cur.fetchone()
8083
except MySQLConnectorError as e:
@@ -107,7 +110,7 @@ def rough_counts_estimator(self, query: str) -> int:
107110

108111
try:
109112
cur.execute(count_query)
110-
logging.info(f"Successfully execute query: {count_query}")
113+
logger.info(f"Successfully execute query: {count_query}")
111114

112115
records = cur.fetchall()
113116

@@ -146,7 +149,7 @@ def record_loader(self, query: str):
146149
try:
147150
cur.execute(query)
148151
conn.commit()
149-
logging.info(f"Successfully execute query: {query}")
152+
logger.info(f"Successfully execute query: {query}")
150153
except MySQLConnectorError as e:
151154
conn.rollback()
152155
raise MySQLClientError(f"Failed to execute query: {query}, Raise exception: {e}")

beam_mysql/connector/io.py

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -115,7 +115,13 @@ def process(self, element: Dict, *args, **kwargs):
115115
values.append(value)
116116

117117
column_str = ", ".join(columns)
118-
value_str = ", ".join([f"{value}" if isinstance(value, (int, float)) else f"'{value}'" for value in values])
118+
value_str = ", ".join(
119+
[
120+
f"{'NULL' if value is None else value}" if isinstance(value, (type(None), int, float)) else f"'{value}'"
121+
for value in values
122+
]
123+
)
124+
119125
query = f"INSERT INTO {self._config['database']}.{self._table}({column_str}) VALUES({value_str});"
120126

121127
self._queries.append(query)

scripts/01_create_test_table.sql

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ CREATE TABLE IF NOT EXISTS `tests`
55
`id` INT(20) AUTO_INCREMENT,
66
`name` VARCHAR(20) NOT NULL,
77
`date` DATE NOT NULL,
8+
`memo` VARCHAR(50),
89
PRIMARY KEY (`id`, `date`)
910
) ENGINE=InnoDB DEFAULT CHARSET=utf8 COLLATE=utf8_bin
1011
PARTITION BY RANGE COLUMNS(date) (

scripts/02_insert_test_data.sql

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
START TRANSACTION;
2-
INSERT INTO tests(name, date) VALUES('test data1', '2020-01-01');
2+
INSERT INTO tests(name, date, memo) VALUES('test data1', '2020-01-01', 'memo1');
33
INSERT INTO tests(name, date) VALUES('test data2', '2020-02-02');
4-
INSERT INTO tests(name, date) VALUES('test data3', '2020-03-03');
4+
INSERT INTO tests(name, date, memo) VALUES('test data3', '2020-03-03', 'memo3');
55
INSERT INTO tests(name, date) VALUES('test data4', '2020-04-04');
66
INSERT INTO tests(name, date) VALUES('test data5', '2020-05-05');
77
COMMIT;

setup.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,7 @@ def get_version():
4848
PACKAGE_DESCRIPTION = "MySQL I/O Connector of Apache Beam"
4949
PACKAGE_URL = "https://github.com/esaki01/beam-mysql-connector"
5050
PACKAGE_DOWNLOAD_URL = "https://pypi.python.org/pypi/beam-mysql-connector"
51-
PACKAGE_AUTHOR = "esaki01"
51+
PACKAGE_AUTHOR = "esakik"
5252
PACKAGE_EMAIL = "esaki1011@gmail.com"
5353
PACKAGE_KEYWORDS = "apache beam mysql connector"
5454
PACKAGE_LONG_DESCRIPTION = README
@@ -77,6 +77,7 @@ def get_version():
7777
"Programming Language :: Python :: 3.6",
7878
"Programming Language :: Python :: 3.7",
7979
"Programming Language :: Python :: 3.8",
80+
"Programming Language :: Python :: 3.9",
8081
],
8182
license="MIT",
8283
keywords=PACKAGE_KEYWORDS,

tests/test_read_records_pipeline.py

Lines changed: 17 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -13,11 +13,11 @@
1313
class TestReadRecordsPipeline(TestBase):
1414
def test_pipeline_no_splitter(self):
1515
expected = [
16-
{"id": 1, "name": "test data1", "date": date(2020, 1, 1)},
17-
{"id": 2, "name": "test data2", "date": date(2020, 2, 2)},
18-
{"id": 3, "name": "test data3", "date": date(2020, 3, 3)},
19-
{"id": 4, "name": "test data4", "date": date(2020, 4, 4)},
20-
{"id": 5, "name": "test data5", "date": date(2020, 5, 5)},
16+
{"id": 1, "name": "test data1", "date": date(2020, 1, 1), "memo": "memo1"},
17+
{"id": 2, "name": "test data2", "date": date(2020, 2, 2), "memo": None},
18+
{"id": 3, "name": "test data3", "date": date(2020, 3, 3), "memo": "memo3"},
19+
{"id": 4, "name": "test data4", "date": date(2020, 4, 4), "memo": None},
20+
{"id": 5, "name": "test data5", "date": date(2020, 5, 5), "memo": None},
2121
]
2222

2323
with TestPipeline() as p:
@@ -38,11 +38,11 @@ def test_pipeline_no_splitter(self):
3838

3939
def test_pipeline_limit_offset_splitter(self):
4040
expected = [
41-
{"id": 1, "name": "test data1", "date": date(2020, 1, 1)},
42-
{"id": 2, "name": "test data2", "date": date(2020, 2, 2)},
43-
{"id": 3, "name": "test data3", "date": date(2020, 3, 3)},
44-
{"id": 4, "name": "test data4", "date": date(2020, 4, 4)},
45-
{"id": 5, "name": "test data5", "date": date(2020, 5, 5)},
41+
{"id": 1, "name": "test data1", "date": date(2020, 1, 1), "memo": "memo1"},
42+
{"id": 2, "name": "test data2", "date": date(2020, 2, 2), "memo": None},
43+
{"id": 3, "name": "test data3", "date": date(2020, 3, 3), "memo": "memo3"},
44+
{"id": 4, "name": "test data4", "date": date(2020, 4, 4), "memo": None},
45+
{"id": 5, "name": "test data5", "date": date(2020, 5, 5), "memo": None},
4646
]
4747

4848
with TestPipeline() as p:
@@ -63,8 +63,8 @@ def test_pipeline_limit_offset_splitter(self):
6363

6464
def test_pipeline_ids_splitter(self):
6565
expected = [
66-
{"id": 1, "name": "test data1", "date": date(2020, 1, 1)},
67-
{"id": 2, "name": "test data2", "date": date(2020, 2, 2)},
66+
{"id": 1, "name": "test data1", "date": date(2020, 1, 1), "memo": "memo1"},
67+
{"id": 2, "name": "test data2", "date": date(2020, 2, 2), "memo": None},
6868
]
6969

7070
with TestPipeline() as p:
@@ -85,9 +85,9 @@ def test_pipeline_ids_splitter(self):
8585

8686
def test_pipeline_date_splitter(self):
8787
expected = [
88-
{"id": 1, "name": "test data1", "date": date(2020, 1, 1)},
89-
{"id": 2, "name": "test data2", "date": date(2020, 2, 2)},
90-
{"id": 3, "name": "test data3", "date": date(2020, 3, 3)},
88+
{"id": 1, "name": "test data1", "date": date(2020, 1, 1), "memo": "memo1"},
89+
{"id": 2, "name": "test data2", "date": date(2020, 2, 2), "memo": None},
90+
{"id": 3, "name": "test data3", "date": date(2020, 3, 3), "memo": "memo3"},
9191
]
9292

9393
with TestPipeline() as p:
@@ -108,8 +108,8 @@ def test_pipeline_date_splitter(self):
108108

109109
def test_pipeline_partitions_splitter(self):
110110
expected = [
111-
{"id": 2, "name": "test data2", "date": date(2020, 2, 2)},
112-
{"id": 3, "name": "test data3", "date": date(2020, 3, 3)},
111+
{"id": 2, "name": "test data2", "date": date(2020, 2, 2), "memo": None},
112+
{"id": 3, "name": "test data3", "date": date(2020, 3, 3), "memo": "memo3"},
113113
]
114114

115115
with TestPipeline() as p:

tests/test_write_records_pipeline.py

Lines changed: 7 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -37,12 +37,12 @@ def tearDown(self):
3737

3838
def test_pipeline(self):
3939
expected = [
40-
{"id": 1, "name": "test data1", "date": datetime.date(2020, 1, 1)},
41-
{"id": 2, "name": "test data2", "date": datetime.date(2020, 2, 2)},
42-
{"id": 3, "name": "test data3", "date": datetime.date(2020, 3, 3)},
43-
{"id": 4, "name": "test data4", "date": datetime.date(2020, 4, 4)},
44-
{"id": 5, "name": "test data5", "date": datetime.date(2020, 5, 5)},
45-
{"id": 6, "name": "test data6", "date": datetime.date(2020, 6, 6)},
40+
{"id": 1, "name": "test data1", "date": datetime.date(2020, 1, 1), "memo": "memo1"},
41+
{"id": 2, "name": "test data2", "date": datetime.date(2020, 2, 2), "memo": None},
42+
{"id": 3, "name": "test data3", "date": datetime.date(2020, 3, 3), "memo": "memo3"},
43+
{"id": 4, "name": "test data4", "date": datetime.date(2020, 4, 4), "memo": None},
44+
{"id": 5, "name": "test data5", "date": datetime.date(2020, 5, 5), "memo": None},
45+
{"id": 6, "name": "test data6", "date": datetime.date(2020, 6, 6), "memo": None},
4646
]
4747

4848
with TestPipeline() as p:
@@ -57,10 +57,7 @@ def test_pipeline(self):
5757
batch_size=BATCH_SIZE,
5858
)
5959

60-
(p
61-
| beam.Create([{"id": 6, "name": "test data6", "date": "2020-06-06"}])
62-
| write_to_mysql
63-
)
60+
(p | beam.Create([{"id": 6, "name": "test data6", "date": "2020-06-06", "memo": None}]) | write_to_mysql)
6461

6562
cur = self.conn.cursor(dictionary=True)
6663
cur.execute(f"SELECT * FROM {DATABASE}.{TABLE}")

0 commit comments

Comments
 (0)