· 8 years ago · May 14, 2018, 01:56 AM
1#!/usr/bin/env python2.6
2# -*- coding: utf-8 -*-
3
4"""
5Optimization
6"""
7import psyco
8psyco.full()
9
10
11"""
12Loading standard library
13"""
14import sys
15import os.path
16import logging
17import logging.config
18import socket
19import urllib2
20import threading
21import datetime
22import Queue as queue
23import ConfigParser as configparser
24
25from optparse import OptionParser
26from urllib2 import Request, URLError, HTTPError, urlopen
27
28
29"""
30Loading external libraries
31"""
32from sqlalchemy import create_engine
33from sqlalchemy import MetaData
34from sqlalchemy import Column
35from sqlalchemy import Integer
36from sqlalchemy import Binary
37from sqlalchemy import String
38from sqlalchemy import DateTime
39from sqlalchemy import ForeignKey
40from sqlalchemy import MetaData
41from sqlalchemy import Table
42from sqlalchemy import ForeignKeyConstraint
43
44from sqlalchemy.sql import exists
45from sqlalchemy.orm import mapper
46from sqlalchemy.orm import sessionmaker
47
48"""
49Config
50"""
51#General
52config = configparser.ConfigParser()
53config.read('%s/etc/pcrawl.conf' % os.path.dirname(__file__))
54base = config.get('main', 'base')
55
56#Logging
57logging.config.fileConfig('%s/%s' % (base, config.get('logging', 'config')))
58log = logging.getLogger('pcrawl')
59
60#Socket
61socket.setdefaulttimeout(config.getint('client', 'timeout'))
62
63#Command line options
64getopt = OptionParser(usage="%prog [[-d IMPORT_DIR] | [-c] | [-a]]", version="%prog 0.1")
65getopt.add_option('-c', '--crawl', action='store_false', default=True,
66 help='Crawl mode')
67getopt.add_option('-a', '--analyze', action='store_true', default=False,
68 help='Analyze mode')
69getopt.add_option('-d', '--import-dir', default=None, help='Import directory')
70
71(options, args) = getopt.parse_args()
72
73#Database
74metadata = MetaData()
75pcrawl_urls = Table('pcrawl_urls', metadata,
76 Column('url', String, primary_key=True)
77 )
78
79pcrawl_pages = Table('pcrawl_pages', metadata,
80 Column('url', String, ForeignKey('pcrawl_urls.url'), primary_key=True),
81 Column('ts', DateTime(timezone=True), primary_key=True,
82 default=datetime.datetime.now),
83 Column('content', Binary, nullable=False)
84 )
85
86pcrawl_features = Table('pcrawl_features', metadata,
87 Column('url', String),
88 Column('ts', DateTime(timezone=True), primary_key=True),
89 Column('name', String),
90 Column('value', Integer, nullable=False),
91 ForeignKeyConstraint(['url', 'ts'],
92 ['pcrawl_pages.url', 'pcrawl_pages.ts'])
93 )
94
95
96"""
97Classes
98"""
99class Url(object):
100 """Implements a web url."""
101
102 def __init__(self, url):
103 self.url = url
104
105 def __repr__(self):
106 return "<Url('%s')>" % self.url
107
108
109class Page(object):
110 """Implements a web page."""
111
112 def __init__(self, content, url, ts=None):
113 if ts is not None:
114 self.ts = ts
115 self.content = content
116 self.url = url
117
118 def __repr__(self):
119 return "<Page('%s', '%s', '%d')>" % (str(self.ts), self.url,
120 len(self.content))
121
122
123class Feature(object):
124 """Implements a feature."""
125
126 def __init__(self, id, url, name, value):
127 self.id = id
128 self.url = url
129 self.name = name
130 self.value = value
131
132 def __repr__(self):
133 return "<Feature('%d', '%s', '%s', '%d')>" % (self.id, self.url,
134 self.name, self.value)
135
136
137class Crawler(threading.Thread):
138 """
139 Implements a crawler thread.
140
141 Consumes urls out of a syncronized queue, crawls and pushes them
142 into an ouptut syncronized queue.
143 """
144 def __init__(self, url_queue, out_queue):
145 """Inintialize variables and calls parent constructor."""
146 threading.Thread.__init__(self)
147 self.url_queue = url_queue
148 self.out_queue = out_queue
149
150 def fetch(self, url):
151 """Fetches a web page and returns a PageSnapshot object."""
152 hdr = { 'User-Agent' : config.get('client', 'ua') }
153 req = Request(url, headers=hdr)
154 try:
155 res = urlopen(req)
156 except HTTPError, e:
157 log.warning('HTTP error opening %s. Reason: %s' % (url, e.code))
158 return None
159 except URLError, e:
160 log.warning('Error opening %s. Reason: %s' % (url, e.reason))
161 return None
162 log.debug('Crawler %d fetched %s' % (self.ident, url) )
163 return Page(res.read(), url)
164
165 def run(self):
166 """Does the job."""
167 log.debug('Crawler %d started' % self.ident)
168 while True:
169 """Reading"""
170 url = self.url_queue.get()
171 self.url_queue.task_done()
172
173 """Fetching"""
174 page = self.fetch(url)
175 if page is not None:
176 self.out_queue.put(page)
177
178
179class Analyzer(threading.Thread):
180 """
181 Calculate features and caches them into the db.
182 """
183 def __init__(self, page_queue, out_queue):
184 threading.Thread.__init__(self)
185 self.page_queue = page_queue
186 self.out_queue = out_queue
187
188 def calc(self, content):
189 return {'forms': content.count('<form '),
190 'links': content.count('<a '),
191 'inputs': content.count('<input ')}
192
193 def run(self):
194 page = self.page_queue.get()
195 counts = self.calc(page.content)
196 for name, value in counts:
197 self.out_queue.put(Feature(page.id, page.url, name, value))
198 self.page_queue.task_done()
199
200
201class Storer(threading.Thread):
202 """
203 Implements a storing thread.
204
205 Consumes objects out of a syncronized queue and pushes them into a
206 transactional db.
207 """
208 def __init__(self, in_queue, db, flush_every_n=100):
209 """Inintialize variables and calls parent constructor."""
210 threading.Thread.__init__(self)
211 self.db = db
212 self.in_queue = in_queue
213 self.flush_every_n = flush_every_n
214 self.counter = 0
215
216 def run(self):
217 """Does the job."""
218 log.debug('Storer %d started (flush every %d objects)' % (self.ident,
219 self.flush_every_n))
220 while True:
221 obj = self.in_queue.get()
222 if self.db is None:
223 log.info('Object: %s' % obj)
224 else:
225 self.db.add(obj)
226 log.debug('Storing %s' % obj)
227 self.in_queue.task_done()
228
229
230def main():
231 """
232 The main
233 """
234
235 """Starting up"""
236 log.info('Cralwer is starting')
237
238 """Create the queues"""
239 tmp_queue = queue.Queue()
240
241 """Initialize the database"""
242 log.info('Initializing database')
243 engine = create_engine(
244 '%s://%s:%s@%s:%s/%s' % (config.get('db', 'driver'),
245 config.get('db', 'user'),
246 config.get('db', 'password'),
247 config.get('db', 'host'),
248 config.get('db', 'port'),
249 config.get('db', 'database')),
250 echo=config.get('db', 'echo'),
251 encoding=config.get('db', 'encoding'))
252 metadata.bind = engine
253 metadata.create_all()
254 mapper(Url, pcrawl_urls)
255 mapper(Page, pcrawl_pages)
256 mapper(Feature, pcrawl_features)
257
258 connection = engine.connect()
259 Session = sessionmaker()
260 db = Session(bind=connection, autoflush=True, autocommit=True)
261
262 """Start the storer thread"""
263 storer = Storer(tmp_queue, db, config.getint('db', 'flush_every_n'))
264 storer.start()
265
266 if options.import_dir is not None:
267 import gzip
268 for site in db.query(Url).\
269 filter(~exists().where(Url.url==Page.url)):
270 for dirpath, dirnames, filenames in os.walk(options.import_dir + site.url[7:]):
271 for file in filenames:
272 t = file[:19].split('-')
273 log.debug('Pushing %s' % dirpath + file)
274 ts = '%s-%s-%s %s:%s:%s%s' % (t[0], t[1], t[2], t[3], t[4], t[5], '-08')
275 tmp_queue.put(Page(content=gzip.open(dirpath + file, 'rb').read(),
276 url=site.url, ts=ts))
277 else:
278 if options.crawl or not options.analyze:
279 url_queue = queue.Queue()
280 for site in db.query(Url):
281 url_queue.put(site.url)
282
283 for i in range(config.getint('client', 'threads')):
284 crawler = Crawler(url_queue, tmp_queue)
285 crawler.start()
286 url_queue.join()
287 else:
288 page_queue = queue.Queue()
289 for page in db.query(Page).\
290 filter(~exists().where(Page.id==Feature.id)):
291 page_queue.put(page)
292 for i in range(config.getint('client', 'threads')):
293 analyzer = Analyzer(page_queue, tmp_queue)
294 analyzer.start()
295 page_queue.join()
296 tmp_queue.join()
297
298 """Terminating"""
299 db.close()
300 connection.close()
301 log.info('Crawler is terminating')
302
303 return 1
304
305
306if __name__ == '__main__':
307 sys.exit(main())