From cdbcbac49d256b284680d6296a8149adedb2935a Mon Sep 17 00:00:00 2001 From: simplypower-bbj Date: Sat, 17 Dec 2022 00:47:51 +0100 Subject: [PATCH] added database interaction class --- database.ipynb | 75 +++++++++++++++++++++++++ database.py | 150 +++++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 225 insertions(+) create mode 100644 database.ipynb create mode 100644 database.py diff --git a/database.ipynb b/database.ipynb new file mode 100644 index 0000000..ee2be67 --- /dev/null +++ b/database.ipynb @@ -0,0 +1,75 @@ +{ + "cells": [ + { + "cell_type": "code", + "execution_count": 2, + "metadata": {}, + "outputs": [ + { + "data": { + "text/plain": [ + "{'username': '@otwndk',\n", + " 'time': datetime.datetime(2022, 12, 16, 21, 19, 9, tzinfo=tzutc()),\n", + " 'content': 'Sponsoreret: Verdensmestre, europamestre og store danske medaljetagere rammer Odense til et prestigefyldt stævne\\n#odense #fyn #cykling'}" + ] + }, + "execution_count": 2, + "metadata": {}, + "output_type": "execute_result" + } + ], + "source": [ + "from database import Data_table\n", + "from scraper import scrape\n", + "\n", + "topic = 'Odense'\n", + "username_list, timestamp_list, content_list = scrape(topic, limit=3, headless=False)\n", + "\n", + "val_dict = {\n", + " 'username': username_list[0],\n", + " 'time': timestamp_list[0],\n", + " 'content': content_list[0]\n", + "}\n", + "\n", + "print(val_dict)\n", + "\n", + "with Data_table() as dt:\n", + " dt.insert(val_dict)" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [] + } + ], + "metadata": { + "kernelspec": { + "display_name": "venv", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.10.6" + }, + "orig_nbformat": 4, + "vscode": { + "interpreter": { + "hash": "1d532f3642617da00cbdbdc5ef98bd8d4ca226c95ac593ade79d8cbc880a2afe" + } + } + }, + "nbformat": 4, + "nbformat_minor": 2 +} diff --git a/database.py b/database.py new file mode 100644 index 0000000..c18a096 --- /dev/null +++ b/database.py @@ -0,0 +1,150 @@ +#!/Users/brian/Projects/twitter_scraper/venv/bin/python +import logging +import os +import traceback +import urllib.parse + +import pandas as pd +import psycopg2 +from dotenv import load_dotenv +from sqlalchemy import create_engine + + +class Database(object): + def __init__(self): + # generate uri + engine = os.getenv('DATABASE_ENGINE') + username = os.getenv('DATABASE_USERNAME') + password = os.getenv('DATABASE_PASSWORD') + password = urllib.parse.quote_plus(password) # escape special characters + host = os.getenv('DATABASE_HOST') + port = os.getenv('DATABASE_PORT') + database = os.getenv('DATABASE_DATABASE') + self._uri = f'{engine}://{username}:{password}@{host}:{port}/{database}' + # prepare connection to database + self._engine = create_engine(self._uri) + self.con = None + + def __enter__(self): + self.connect() + return self + + def __exit__(self, exc_type, exc_value, tb): + if not exc_type is None: + traceback.print_exception(exc_type, exc_value, tb) + self.disconnect() + + def connect(self): + try: + con = self._engine.connect() + except psycopg2.OperationalError as e: + logging.warning(e) + else: + self.con = con + + def disconnect(self): + self.con.close() + + def execute(self, query, var=None): + assert self.con is not None, 'connection to database could not be established' + try: + if var is None: + self.con.execute(query) + else: + self.con.execute(query, var) + except Exception as e: + print(e) + raise + + def get(self, query): + df = pd.read_sql_query(query, self.con) + return df + +class Schema(Database): + def __init__(self, name): + super(Schema, self).__init__() + self._name = name + + def __enter__(self): + self.connect() + if not self.exists(): + self.create() + return self + + def exists(self): + query = ( + 'SELECT schema_name ' + 'FROM information_schema.schemata ' + 'ORDER BY schema_name ' + ';' + ) + df = self.get(query) + schema_list = df['schema_name'].tolist() + return self._name in schema_list + + def create(self): + query = f'CREATE SCHEMA {self._name};' + self.execute(query) + + def purge(self): + query = f'DROP SCHEMA IF EXISTS {self._name} CASCADE;' + self.execute(query) + self.create() + +class Table(Database): + def __init__(self, name): + super(Table, self).__init__() + self._name = name + + def __enter__(self): + self.connect() + if not self.exists(): + self.create() + return self + + def exists(self): + query = ( + 'SELECT table_name ' + 'FROM information_schema.tables ' + "WHERE table_schema='public' " + 'ORDER BY table_name ' + ) + df = self.get(query) + table_list = df['table_name'].tolist() + return self._name in table_list + + def create(self): + # placeholder for childrens' methods + pass + + def purge(self): + query = f'DROP TABLE IF EXISTS {self._name} CASCADE;' + self.execute(query) + self.create() + +class Data_table(Table): + def __init__(self): + super(Data_table, self).__init__('data_table') + + def create(self): + query = ( + f'CREATE TABLE {self._name} ' + '(' + 'id SERIAL PRIMARY KEY, ' + 'username VARCHAR(15) NOT NULL, ' + 'time TIMESTAMPTZ NOT NULL, ' + 'content TEXT NOT NULL, ' + 'UNIQUE (username, time, content) ' + ')' + ';' + ) + self.execute(query) + + def insert(self, val_dict): + query = ( + f'INSERT INTO {self._name} (username, time, content) ' + 'VALUES (%(username)s, %(time)s, %(content)s) ' + 'ON CONFLICT (username, time, content) DO NOTHING ' + ';' + ) + self.execute(query, val_dict) \ No newline at end of file