作者BlgAtlfans (BLG_Eric)
看板Python
标题[问题] django celery问题
时间Thu Oct 27 15:37:47 2016
各位大大好
最近在用celery处理
csv,xlsx档案写入postgresql的功能
但是有一些问题想请教
P.S下面的程式码没加入celery时都可以正常执行(大档案例外)
1. csv档写入时celery debug有以下错误
http://imgur.com/8Fg3dmJ
http://imgur.com/DSpZU1s
附上task.py程式码
# -*- coding: utf-8 -*-
from django.shortcuts import render_to_response
from django.template import RequestContext
from django.http import HttpResponseRedirect
from django.core.urlresolvers import reverse
from django.contrib import messages
from django.conf import settings
from django.db import connection
from django.views.decorators.csrf import csrf_exempt
from celery import Celery
from celery import task
import json
import csv
import sys
import random
import psycopg2
import xlrd
import openpyxl as pyxl
from .models import Document
from .forms import DocumentForm
app = Celery(
'tasks',
broker='amqp://guest:guest@localhost:5672//',
backend='rpc://'
)
CELERY_RESULT_BACKEND = 'rpc://'
CELERY_RESULT_PERSISTENT = False
@app.task()
def csvwritein(doc):# Transform csv to table
doc = doc
conn = psycopg2.connect("dbname='apidb' user='api' host='localhost'
password='eric40502' port='5432'")
readcur = conn.cursor()
readcur.execute("select exists(select * from
information_schema.tables where table_name='%s')" % doc.tablename) # check if
same file is already in database
check = readcur.fetchone()[0]
try:
fr = open(doc.path,encoding = 'utf-8-sig')
dr.delay(fr,doc,check)
fr.close()
except Exception as e:
fr = open(doc.path,encoding = 'big5')
dr.delay(fr,doc,check)
fr.close()
conn.commit()
readcur.close()
@app.task()
def dr(fr,doc,check): # make datareader as function to keep code 'dry'
csvt = 0 #count csv reader loop time
row_id = 1 # used for following id field
conn = psycopg2.connect("dbname='apidb' user='api' host='localhost'
password='eric40502' port='5432'")
maincur = conn.cursor()
writecur = conn.cursor()
datareader = csv.reader(fr, delimiter=',')
for row in datareader:
if csvt == 0: # first time in loop(create field) and check no
same file exists
if check == True:
app =
''.join([random.SystemRandom().choice('abcdefghijklmnopqrstuvwxyz0123456789')
for i in range(6)])
tname = '%s-%s' % (doc.tablename,app
tablename = '"%s-%s"' % (doc.tablename,app)
doc.tablename = tname
doc.save()
else:
tablename = '"%s"' % doc.tablename
maincur.execute("CREATE TABLE %s (id SERIAL PRIMARY
KEY);" % tablename)
row_count = sum(1 for line in datareader)
col_count = len(row)
frow = row
for i in range(0,col_count,1):
row[i] = '"%s"' % row[i] # change number to
string
maincur.execute("ALTER TABLE %s ADD %s
CITEXT;" % (tablename,row[i]))
csvt = csvt+1
fr.seek(0)
next(datareader)
elif csvt > 0: # not first time(insert data) and check no
same file exists
for j in range(0,col_count,1):
if j == 0:
writecur.execute("INSERT INTO %s (%s)
VALUES ('%s');" % (tablename,frow[j],row[j]))
else:
writecur.execute("UPDATE %s SET %s =
'%s' WHERE id = '%d';" %(tablename,frow[j],row[j],row_id))
csvt = csvt+1
row_id = row_id+1
else:
break
conn.commit()
maincur.close()
writecur.close()
conn.close()
csvt = 0
doc = Document.objects.all()
呼叫时是用csvwritein.delay(doc)
2. xlsx 档案(13万笔资料)写入时 worker卡了两三分钟後跑出以下错误
http://imgur.com/rC1uun0
以下是task.py xlsx写入函数
@app.task()
def xlsxwritein(doc): # write into database for file type xlsx
xlsxt = 0
conn = psycopg2.connect("dbname='apidb' user='api' host='localhost'
password='eric40502' port='5432'")
maincur = conn.cursor()
readcur = conn.cursor()
writecur = conn.cursor()
readcur.execute("select exists(select * from
information_schema.tables where table_name='%s')" % doc.tablename) # check if
same file is already in database
check = readcur.fetchone()[0]
row_id = 1 # used for following id field
wb = pyxl.load_workbook(doc.path)
sheetnames = wb.get_sheet_names()
ws = wb.get_sheet_by_name(sheetnames[0])
for rown in range(ws.get_highest_row()):
if xlsxt == 0:
if check == True:
app =
''.join([random.SystemRandom().choice('abcdefghijklmnopqrstuvwxyz0123456789')
for i in range(6)])
tname = '%s-%s' % (doc.tablename,app)
tablename = '"%s-%s"' % (doc.tablename,app)
doc.tablename = tname
doc.save()
else:
tablename = '"%s"' % doc.tablename
field = [ws.cell(row=1,column=col_index).value for
col_index in range(1,ws.get_highest_column()+1)]
maincur.execute("CREATE TABLE %s (id SERIAL PRIMARY
KEY);" % tablename)
for coln in range(ws.get_highest_column()):
field[coln] = '"%s"' % field[coln] # change
number to string
if field[coln] == 'ID':
field[coln] = 'original_id'
maincur.execute("ALTER TABLE %s ADD %s
CITEXT;" % (tablename,field[coln]))
xlsxt = xlsxt+1
elif xlsxt > 0 and check == False: # not first time(insert
data) and check no same file exists
for coln in range(ws.get_highest_column()):
if coln == 0:
writecur.execute("INSERT INTO %s (%s)
VALUES ('%s');"
%(tablename,field[coln],str(ws.cell(row=rown,column=coln+1).value)))
else:
writecur.execute("UPDATE %s SET %s =
'%s' WHERE id = '%d';"
%(tablename,field[coln],str(ws.cell(row=rown+1,column=coln+1).value),row_id))
xlsxt = xlsxt+1
row_id = row_id+1
else:
break
conn.commit()
maincur.close()
readcur.close()
writecur.close()
conn.close()
xlsxt = 0
--
※ 发信站: 批踢踢实业坊(ptt.cc), 来自: 114.32.19.185
※ 文章网址: https://webptt.com/cn.aspx?n=bbs/Python/M.1477553869.A.5F0.html
1F:→ kenduest: 好像是被作业系统 kernel 踢出去了? 10/27 16:12
2F:→ kenduest: 比方吃太多记忆体等,被 linux OOM killer 处理掉 10/27 16:12
3F:→ BlgAtlfans: 那应该要怎麽样处理 10/27 16:52
4F:→ kenduest: 你先独立把那个处理task写成独立档案单独终端跑看看 10/27 18:33
5F:→ kenduest: 後续用 free 与 top 看一下记忆体使用情况 10/27 18:33
6F:→ kenduest: 或许是实际那个 server 本来记忆体就不多所以就爆掉了 10/27 18:34
7F:→ kenduest: 题外话你的程式码贴这边很乱很难看清楚 10/27 18:40
8F:→ kenduest: 另外建议请用 4 个空白代替 tab, 建议这样在 python 上 10/27 18:41
9F:→ uranusjr: 先试试看 DEBUG = False, 这两个记忆体用量差很多 10/27 20:24
10F:→ uranusjr: 第一个问题要看你 doc 到底是什麽 10/27 20:24
11F:→ BlgAtlfans: 感谢各位回答 我的doc是个django的model 10/27 21:26
12F:→ BlgAtlfans: 内容是上传档案的一些资料 10/27 21:27
13F:→ BlgAtlfans: 像是tablename path id...之类的 10/27 21:28
14F:→ BlgAtlfans: 这里主要是用来传递tablename来做为建table的依据 10/27 21:30
15F:→ BlgAtlfans: 多问一个 一般来说写入一个13万行的资料需要很多记忆 10/27 21:32
16F:→ BlgAtlfans: 体吗? 10/27 21:32