Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Submit feedback
Contribute to GitLab
Sign in
Toggle navigation
W
webmagic
Project
Project
Details
Activity
Releases
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
沈俊林
webmagic
Commits
0f2c5b57
Commit
0f2c5b57
authored
Aug 11, 2013
by
yihua.huang
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
update redisscheduler
parent
787b9529
Changes
10
Show whitespace changes
Inline
Side-by-side
Showing
10 changed files
with
157 additions
and
95 deletions
+157
-95
Page.java
webmagic-core/src/main/java/us/codecraft/webmagic/Page.java
+11
-0
Request.java
...gic-core/src/main/java/us/codecraft/webmagic/Request.java
+9
-0
Spider.java
...agic-core/src/main/java/us/codecraft/webmagic/Spider.java
+19
-19
HttpClientDownloader.java
...s/codecraft/webmagic/downloader/HttpClientDownloader.java
+11
-7
FilePipeline.java
...ain/java/us/codecraft/webmagic/pipeline/FilePipeline.java
+6
-15
FilePersistentBase.java
.../java/us/codecraft/webmagic/utils/FilePersistentBase.java
+51
-0
OOSpider.java
...n/src/main/java/us/codecraft/webmagic/model/OOSpider.java
+5
-0
JsonFilePageModelPipeline.java
...odecraft/webmagic/pipeline/JsonFilePageModelPipeline.java
+7
-16
JsonFilePipeline.java
...java/us/codecraft/webmagic/pipeline/JsonFilePipeline.java
+7
-15
RedisScheduler.java
.../java/us/codecraft/webmagic/scheduler/RedisScheduler.java
+31
-23
No files found.
webmagic-core/src/main/java/us/codecraft/webmagic/Page.java
View file @
0f2c5b57
...
@@ -148,4 +148,15 @@ public class Page {
...
@@ -148,4 +148,15 @@ public class Page {
public
ResultItems
getResultItems
()
{
public
ResultItems
getResultItems
()
{
return
resultItems
;
return
resultItems
;
}
}
@Override
public
String
toString
()
{
return
"Page{"
+
"request="
+
request
+
", resultItems="
+
resultItems
+
", html="
+
html
+
", url="
+
url
+
", targetRequests="
+
targetRequests
+
'}'
;
}
}
}
webmagic-core/src/main/java/us/codecraft/webmagic/Request.java
View file @
0f2c5b57
...
@@ -113,4 +113,13 @@ public class Request implements Serializable {
...
@@ -113,4 +113,13 @@ public class Request implements Serializable {
public
void
setUrl
(
String
url
)
{
public
void
setUrl
(
String
url
)
{
this
.
url
=
url
;
this
.
url
=
url
;
}
}
@Override
public
String
toString
()
{
return
"Request{"
+
"url='"
+
url
+
'\''
+
", extras="
+
extras
+
", priority="
+
priority
+
'}'
;
}
}
}
webmagic-core/src/main/java/us/codecraft/webmagic/Spider.java
View file @
0f2c5b57
...
@@ -40,33 +40,33 @@ import java.util.concurrent.atomic.AtomicInteger;
...
@@ -40,33 +40,33 @@ import java.util.concurrent.atomic.AtomicInteger;
*/
*/
public
class
Spider
implements
Runnable
,
Task
{
public
class
Spider
implements
Runnable
,
Task
{
pr
ivate
Downloader
downloader
;
pr
otected
Downloader
downloader
;
pr
ivate
List
<
Pipeline
>
pipelines
=
new
ArrayList
<
Pipeline
>();
pr
otected
List
<
Pipeline
>
pipelines
=
new
ArrayList
<
Pipeline
>();
pr
ivate
PageProcessor
pageProcessor
;
pr
otected
PageProcessor
pageProcessor
;
pr
ivate
List
<
String
>
startUrls
;
pr
otected
List
<
String
>
startUrls
;
pr
ivate
Site
site
;
pr
otected
Site
site
;
pr
ivate
String
uuid
;
pr
otected
String
uuid
;
pr
ivate
Scheduler
scheduler
=
new
QueueScheduler
();
pr
otected
Scheduler
scheduler
=
new
QueueScheduler
();
pr
ivate
Logger
logger
=
Logger
.
getLogger
(
getClass
());
pr
otected
Logger
logger
=
Logger
.
getLogger
(
getClass
());
pr
ivate
ExecutorService
executorService
;
pr
otected
ExecutorService
executorService
;
pr
ivate
int
threadNum
=
1
;
pr
otected
int
threadNum
=
1
;
pr
ivate
AtomicInteger
stat
=
new
AtomicInteger
(
STAT_INIT
);
pr
otected
AtomicInteger
stat
=
new
AtomicInteger
(
STAT_INIT
);
pr
ivate
final
static
int
STAT_INIT
=
0
;
pr
otected
final
static
int
STAT_INIT
=
0
;
pr
ivate
final
static
int
STAT_RUNNING
=
1
;
pr
otected
final
static
int
STAT_RUNNING
=
1
;
pr
ivate
final
static
int
STAT_STOPPED
=
2
;
pr
otected
final
static
int
STAT_STOPPED
=
2
;
/**
/**
* 使用已定义的抽取规则新建一个Spider。
* 使用已定义的抽取规则新建一个Spider。
...
@@ -206,7 +206,7 @@ public class Spider implements Runnable, Task {
...
@@ -206,7 +206,7 @@ public class Spider implements Runnable, Task {
destroy
();
destroy
();
}
}
pr
ivate
void
destroy
()
{
pr
otected
void
destroy
()
{
destroyEach
(
downloader
);
destroyEach
(
downloader
);
destroyEach
(
pageProcessor
);
destroyEach
(
pageProcessor
);
for
(
Pipeline
pipeline
:
pipelines
)
{
for
(
Pipeline
pipeline
:
pipelines
)
{
...
@@ -233,7 +233,7 @@ public class Spider implements Runnable, Task {
...
@@ -233,7 +233,7 @@ public class Spider implements Runnable, Task {
}
}
}
}
pr
ivate
void
processRequest
(
Request
request
)
{
pr
otected
void
processRequest
(
Request
request
)
{
Page
page
=
downloader
.
download
(
request
,
this
);
Page
page
=
downloader
.
download
(
request
,
this
);
if
(
page
==
null
)
{
if
(
page
==
null
)
{
sleep
(
site
.
getSleepTime
());
sleep
(
site
.
getSleepTime
());
...
@@ -249,7 +249,7 @@ public class Spider implements Runnable, Task {
...
@@ -249,7 +249,7 @@ public class Spider implements Runnable, Task {
sleep
(
site
.
getSleepTime
());
sleep
(
site
.
getSleepTime
());
}
}
pr
ivate
void
sleep
(
int
time
)
{
pr
otected
void
sleep
(
int
time
)
{
try
{
try
{
Thread
.
sleep
(
time
);
Thread
.
sleep
(
time
);
}
catch
(
InterruptedException
e
)
{
}
catch
(
InterruptedException
e
)
{
...
@@ -257,7 +257,7 @@ public class Spider implements Runnable, Task {
...
@@ -257,7 +257,7 @@ public class Spider implements Runnable, Task {
}
}
}
}
pr
ivate
void
addRequest
(
Page
page
)
{
pr
otected
void
addRequest
(
Page
page
)
{
if
(
CollectionUtils
.
isNotEmpty
(
page
.
getTargetRequests
()))
{
if
(
CollectionUtils
.
isNotEmpty
(
page
.
getTargetRequests
()))
{
for
(
Request
request
:
page
.
getTargetRequests
())
{
for
(
Request
request
:
page
.
getTargetRequests
())
{
scheduler
.
push
(
request
,
this
);
scheduler
.
push
(
request
,
this
);
...
@@ -265,7 +265,7 @@ public class Spider implements Runnable, Task {
...
@@ -265,7 +265,7 @@ public class Spider implements Runnable, Task {
}
}
}
}
pr
ivate
void
checkIfNotRunning
()
{
pr
otected
void
checkIfNotRunning
()
{
if
(!
stat
.
compareAndSet
(
STAT_INIT
,
STAT_INIT
))
{
if
(!
stat
.
compareAndSet
(
STAT_INIT
,
STAT_INIT
))
{
throw
new
IllegalStateException
(
"Spider is already running!"
);
throw
new
IllegalStateException
(
"Spider is already running!"
);
}
}
...
...
webmagic-core/src/main/java/us/codecraft/webmagic/downloader/HttpClientDownloader.java
View file @
0f2c5b57
...
@@ -66,13 +66,7 @@ public class HttpClientDownloader implements Downloader {
...
@@ -66,13 +66,7 @@ public class HttpClientDownloader implements Downloader {
}
}
//
//
handleGzip
(
httpResponse
);
handleGzip
(
httpResponse
);
String
content
=
IOUtils
.
toString
(
httpResponse
.
getEntity
().
getContent
(),
return
handleResponse
(
request
,
charset
,
httpResponse
,
task
);
charset
);
Page
page
=
new
Page
();
page
.
setHtml
(
new
Html
(
UrlUtils
.
fixAllRelativeHrefs
(
content
,
request
.
getUrl
())));
page
.
setUrl
(
new
PlainText
(
request
.
getUrl
()));
page
.
setRequest
(
request
);
return
page
;
}
else
{
}
else
{
logger
.
warn
(
"code error "
+
statusCode
+
"\t"
+
request
.
getUrl
());
logger
.
warn
(
"code error "
+
statusCode
+
"\t"
+
request
.
getUrl
());
}
}
...
@@ -82,6 +76,16 @@ public class HttpClientDownloader implements Downloader {
...
@@ -82,6 +76,16 @@ public class HttpClientDownloader implements Downloader {
return
null
;
return
null
;
}
}
protected
Page
handleResponse
(
Request
request
,
String
charset
,
HttpResponse
httpResponse
,
Task
task
)
throws
IOException
{
String
content
=
IOUtils
.
toString
(
httpResponse
.
getEntity
().
getContent
(),
charset
);
Page
page
=
new
Page
();
page
.
setHtml
(
new
Html
(
UrlUtils
.
fixAllRelativeHrefs
(
content
,
request
.
getUrl
())));
page
.
setUrl
(
new
PlainText
(
request
.
getUrl
()));
page
.
setRequest
(
request
);
return
page
;
}
@Override
@Override
public
void
setThread
(
int
thread
)
{
public
void
setThread
(
int
thread
)
{
poolSize
=
thread
;
poolSize
=
thread
;
...
...
webmagic-core/src/main/java/us/codecraft/webmagic/pipeline/FilePipeline.java
View file @
0f2c5b57
...
@@ -4,8 +4,8 @@ import org.apache.commons.codec.digest.DigestUtils;
...
@@ -4,8 +4,8 @@ import org.apache.commons.codec.digest.DigestUtils;
import
org.apache.log4j.Logger
;
import
org.apache.log4j.Logger
;
import
us.codecraft.webmagic.ResultItems
;
import
us.codecraft.webmagic.ResultItems
;
import
us.codecraft.webmagic.Task
;
import
us.codecraft.webmagic.Task
;
import
us.codecraft.webmagic.utils.FilePersistentBase
;
import
java.io.File
;
import
java.io.FileWriter
;
import
java.io.FileWriter
;
import
java.io.IOException
;
import
java.io.IOException
;
import
java.io.PrintWriter
;
import
java.io.PrintWriter
;
...
@@ -18,9 +18,7 @@ import java.util.Map;
...
@@ -18,9 +18,7 @@ import java.util.Map;
* Date: 13-4-21
* Date: 13-4-21
* Time: 下午6:28
* Time: 下午6:28
*/
*/
public
class
FilePipeline
implements
Pipeline
{
public
class
FilePipeline
extends
FilePersistentBase
implements
Pipeline
{
private
String
path
=
"/data/webmagic/"
;
private
Logger
logger
=
Logger
.
getLogger
(
getClass
());
private
Logger
logger
=
Logger
.
getLogger
(
getClass
());
...
@@ -28,7 +26,7 @@ public class FilePipeline implements Pipeline {
...
@@ -28,7 +26,7 @@ public class FilePipeline implements Pipeline {
* 新建一个FilePipeline,使用默认保存路径"/data/webmagic/"
* 新建一个FilePipeline,使用默认保存路径"/data/webmagic/"
*/
*/
public
FilePipeline
()
{
public
FilePipeline
()
{
setPath
(
"/data/webmagic/"
);
}
}
/**
/**
...
@@ -37,21 +35,14 @@ public class FilePipeline implements Pipeline {
...
@@ -37,21 +35,14 @@ public class FilePipeline implements Pipeline {
* @param path 文件保存路径
* @param path 文件保存路径
*/
*/
public
FilePipeline
(
String
path
)
{
public
FilePipeline
(
String
path
)
{
if
(!
path
.
endsWith
(
"/"
)&&!
path
.
endsWith
(
"\\"
)){
setPath
(
path
);
path
+=
"/"
;
}
this
.
path
=
path
;
}
}
@Override
@Override
public
void
process
(
ResultItems
resultItems
,
Task
task
)
{
public
void
process
(
ResultItems
resultItems
,
Task
task
)
{
String
path
=
this
.
path
+
"/"
+
task
.
getUUID
()
+
"/"
;
String
path
=
this
.
path
+
PATH_SEPERATOR
+
task
.
getUUID
()
+
PATH_SEPERATOR
;
File
file
=
new
File
(
path
);
if
(!
file
.
exists
())
{
file
.
mkdirs
();
}
try
{
try
{
PrintWriter
printWriter
=
new
PrintWriter
(
new
FileWriter
(
path
+
DigestUtils
.
md5Hex
(
resultItems
.
getRequest
().
getUrl
())
+
".html"
));
PrintWriter
printWriter
=
new
PrintWriter
(
new
FileWriter
(
getFile
(
path
+
DigestUtils
.
md5Hex
(
resultItems
.
getRequest
().
getUrl
())
+
".html"
)
));
printWriter
.
println
(
"url:\t"
+
resultItems
.
getRequest
().
getUrl
());
printWriter
.
println
(
"url:\t"
+
resultItems
.
getRequest
().
getUrl
());
for
(
Map
.
Entry
<
String
,
Object
>
entry
:
resultItems
.
getAll
().
entrySet
())
{
for
(
Map
.
Entry
<
String
,
Object
>
entry
:
resultItems
.
getAll
().
entrySet
())
{
if
(
entry
.
getValue
()
instanceof
Iterable
)
{
if
(
entry
.
getValue
()
instanceof
Iterable
)
{
...
...
webmagic-core/src/main/java/us/codecraft/webmagic/utils/FilePersistentBase.java
0 → 100644
View file @
0f2c5b57
package
us
.
codecraft
.
webmagic
.
utils
;
import
java.io.File
;
/**
* 文件持久化的基础类。<br>
*
* @author code4crafter@gmail.com <br>
* Date: 13-8-11 <br>
* Time: 下午4:21 <br>
*/
public
class
FilePersistentBase
{
protected
String
path
;
public
static
String
PATH_SEPERATOR
=
"/"
;
static
{
String
property
=
System
.
getProperties
().
getProperty
(
"file.separator"
);
if
(
property
!=
null
)
{
PATH_SEPERATOR
=
property
;
}
}
public
void
setPath
(
String
path
)
{
this
.
path
=
path
;
if
(!
path
.
endsWith
(
PATH_SEPERATOR
))
{
path
+=
PATH_SEPERATOR
;
}
}
public
File
getFile
(
String
fullName
)
{
checkAndMakeParentDirecotry
(
fullName
);
return
new
File
(
fullName
);
}
public
void
checkAndMakeParentDirecotry
(
String
fullName
)
{
int
index
=
fullName
.
lastIndexOf
(
PATH_SEPERATOR
);
if
(
index
>
0
)
{
String
path
=
fullName
.
substring
(
0
,
index
);
File
file
=
new
File
(
path
);
if
(!
file
.
exists
())
{
file
.
mkdirs
();
}
}
}
public
String
getPath
()
{
return
path
;
}
}
webmagic-extension/src/main/java/us/codecraft/webmagic/model/OOSpider.java
View file @
0f2c5b57
...
@@ -2,6 +2,7 @@ package us.codecraft.webmagic.model;
...
@@ -2,6 +2,7 @@ package us.codecraft.webmagic.model;
import
us.codecraft.webmagic.Site
;
import
us.codecraft.webmagic.Site
;
import
us.codecraft.webmagic.Spider
;
import
us.codecraft.webmagic.Spider
;
import
us.codecraft.webmagic.processor.PageProcessor
;
/**
/**
* 基于Model的Spider,封装后的入口类。<br>
* 基于Model的Spider,封装后的入口类。<br>
...
@@ -20,6 +21,10 @@ public class OOSpider extends Spider {
...
@@ -20,6 +21,10 @@ public class OOSpider extends Spider {
this
.
modelPageProcessor
=
modelPageProcessor
;
this
.
modelPageProcessor
=
modelPageProcessor
;
}
}
public
OOSpider
(
PageProcessor
pageProcessor
)
{
super
(
pageProcessor
);
}
/**
/**
* 创建一个爬虫。<br>
* 创建一个爬虫。<br>
* @param site
* @param site
...
...
webmagic-extension/src/main/java/us/codecraft/webmagic/pipeline/JsonFilePageModelPipeline.java
View file @
0f2c5b57
...
@@ -7,8 +7,8 @@ import org.apache.log4j.Logger;
...
@@ -7,8 +7,8 @@ import org.apache.log4j.Logger;
import
us.codecraft.webmagic.Task
;
import
us.codecraft.webmagic.Task
;
import
us.codecraft.webmagic.model.HasKey
;
import
us.codecraft.webmagic.model.HasKey
;
import
us.codecraft.webmagic.model.PageModelPipeline
;
import
us.codecraft.webmagic.model.PageModelPipeline
;
import
us.codecraft.webmagic.utils.FilePersistentBase
;
import
java.io.File
;
import
java.io.FileWriter
;
import
java.io.FileWriter
;
import
java.io.IOException
;
import
java.io.IOException
;
import
java.io.PrintWriter
;
import
java.io.PrintWriter
;
...
@@ -21,38 +21,29 @@ import java.io.PrintWriter;
...
@@ -21,38 +21,29 @@ import java.io.PrintWriter;
* Date: 13-4-21
* Date: 13-4-21
* Time: 下午6:28
* Time: 下午6:28
*/
*/
public
class
JsonFilePageModelPipeline
implements
PageModelPipeline
{
public
class
JsonFilePageModelPipeline
extends
FilePersistentBase
implements
PageModelPipeline
{
private
String
path
=
"/data/webmagic/"
;
private
Logger
logger
=
Logger
.
getLogger
(
getClass
());
private
Logger
logger
=
Logger
.
getLogger
(
getClass
());
/**
/**
* 新建一个
File
Pipeline,使用默认保存路径"/data/webmagic/"
* 新建一个
JsonFilePageModel
Pipeline,使用默认保存路径"/data/webmagic/"
*/
*/
public
JsonFilePageModelPipeline
()
{
public
JsonFilePageModelPipeline
()
{
setPath
(
"/data/webmagic/"
);
}
}
/**
/**
* 新建一个
File
Pipeline
* 新建一个
JsonFilePageModel
Pipeline
*
*
* @param path 文件保存路径
* @param path 文件保存路径
*/
*/
public
JsonFilePageModelPipeline
(
String
path
)
{
public
JsonFilePageModelPipeline
(
String
path
)
{
if
(!
path
.
endsWith
(
"/"
)
&&
!
path
.
endsWith
(
"\\"
))
{
setPath
(
path
);
path
+=
"/"
;
}
this
.
path
=
path
;
}
}
@Override
@Override
public
void
process
(
Object
o
,
Task
task
)
{
public
void
process
(
Object
o
,
Task
task
)
{
String
path
=
this
.
path
+
"/"
+
task
.
getUUID
()
+
"/"
;
String
path
=
this
.
path
+
"/"
+
task
.
getUUID
()
+
"/"
;
File
file
=
new
File
(
path
);
if
(!
file
.
exists
())
{
file
.
mkdirs
();
}
try
{
try
{
String
filename
;
String
filename
;
if
(
o
instanceof
HasKey
)
{
if
(
o
instanceof
HasKey
)
{
...
@@ -60,7 +51,7 @@ public class JsonFilePageModelPipeline implements PageModelPipeline {
...
@@ -60,7 +51,7 @@ public class JsonFilePageModelPipeline implements PageModelPipeline {
}
else
{
}
else
{
filename
=
path
+
DigestUtils
.
md5Hex
(
ToStringBuilder
.
reflectionToString
(
o
))
+
".json"
;
filename
=
path
+
DigestUtils
.
md5Hex
(
ToStringBuilder
.
reflectionToString
(
o
))
+
".json"
;
}
}
PrintWriter
printWriter
=
new
PrintWriter
(
new
FileWriter
(
filename
));
PrintWriter
printWriter
=
new
PrintWriter
(
new
FileWriter
(
getFile
(
filename
)
));
printWriter
.
write
(
JSON
.
toJSONString
(
o
));
printWriter
.
write
(
JSON
.
toJSONString
(
o
));
printWriter
.
close
();
printWriter
.
close
();
}
catch
(
IOException
e
)
{
}
catch
(
IOException
e
)
{
...
...
webmagic-extension/src/main/java/us/codecraft/webmagic/pipeline/JsonFilePipeline.java
View file @
0f2c5b57
...
@@ -5,6 +5,7 @@ import org.apache.commons.codec.digest.DigestUtils;
...
@@ -5,6 +5,7 @@ import org.apache.commons.codec.digest.DigestUtils;
import
org.apache.log4j.Logger
;
import
org.apache.log4j.Logger
;
import
us.codecraft.webmagic.ResultItems
;
import
us.codecraft.webmagic.ResultItems
;
import
us.codecraft.webmagic.Task
;
import
us.codecraft.webmagic.Task
;
import
us.codecraft.webmagic.utils.FilePersistentBase
;
import
java.io.File
;
import
java.io.File
;
import
java.io.FileWriter
;
import
java.io.FileWriter
;
...
@@ -18,40 +19,31 @@ import java.io.PrintWriter;
...
@@ -18,40 +19,31 @@ import java.io.PrintWriter;
* Date: 13-4-21
* Date: 13-4-21
* Time: 下午6:28
* Time: 下午6:28
*/
*/
public
class
JsonFilePipeline
implements
Pipeline
{
public
class
JsonFilePipeline
extends
FilePersistentBase
implements
Pipeline
{
private
String
path
=
"/data/webmagic/"
;
private
Logger
logger
=
Logger
.
getLogger
(
getClass
());
private
Logger
logger
=
Logger
.
getLogger
(
getClass
());
/**
/**
* 新建一个FilePipeline,使用默认保存路径"/data/webmagic/"
* 新建一个
Json
FilePipeline,使用默认保存路径"/data/webmagic/"
*/
*/
public
JsonFilePipeline
()
{
public
JsonFilePipeline
()
{
setPath
(
"/data/webmagic"
);
}
}
/**
/**
* 新建一个FilePipeline
* 新建一个
Json
FilePipeline
*
*
* @param path 文件保存路径
* @param path 文件保存路径
*/
*/
public
JsonFilePipeline
(
String
path
)
{
public
JsonFilePipeline
(
String
path
)
{
if
(!
path
.
endsWith
(
"/"
)&&!
path
.
endsWith
(
"\\"
)){
setPath
(
path
);
path
+=
"/"
;
}
this
.
path
=
path
;
}
}
@Override
@Override
public
void
process
(
ResultItems
resultItems
,
Task
task
)
{
public
void
process
(
ResultItems
resultItems
,
Task
task
)
{
String
path
=
this
.
path
+
"/"
+
task
.
getUUID
()
+
"/"
;
String
path
=
this
.
path
+
"/"
+
task
.
getUUID
()
+
"/"
;
File
file
=
new
File
(
path
);
if
(!
file
.
exists
())
{
file
.
mkdirs
();
}
try
{
try
{
PrintWriter
printWriter
=
new
PrintWriter
(
new
FileWriter
(
path
+
DigestUtils
.
md5Hex
(
resultItems
.
getRequest
().
getUrl
())
+
".json"
));
PrintWriter
printWriter
=
new
PrintWriter
(
new
FileWriter
(
new
File
(
path
+
DigestUtils
.
md5Hex
(
resultItems
.
getRequest
().
getUrl
())
+
".json"
)
));
printWriter
.
write
(
JSON
.
toJSONString
(
resultItems
.
getAll
()));
printWriter
.
write
(
JSON
.
toJSONString
(
resultItems
.
getAll
()));
printWriter
.
close
();
printWriter
.
close
();
}
catch
(
IOException
e
)
{
}
catch
(
IOException
e
)
{
...
...
webmagic-extension/src/main/java/us/codecraft/webmagic/scheduler/RedisScheduler.java
View file @
0f2c5b57
...
@@ -32,34 +32,42 @@ public class RedisScheduler implements Scheduler {
...
@@ -32,34 +32,42 @@ public class RedisScheduler implements Scheduler {
@Override
@Override
public
synchronized
void
push
(
Request
request
,
Task
task
)
{
public
synchronized
void
push
(
Request
request
,
Task
task
)
{
Jedis
jedis
=
pool
.
getResource
();
Jedis
jedis
=
pool
.
getResource
();
//使用SortedSet进行url去重
try
{
if
(
jedis
.
zrank
(
SET_PREFIX
+
task
.
getUUID
(),
request
.
getUrl
())
==
null
)
{
//使用Set进行url去重
if
(!
jedis
.
sismember
(
SET_PREFIX
+
task
.
getUUID
(),
request
.
getUrl
()))
{
//使用List保存队列
//使用List保存队列
jedis
.
rpush
(
QUEUE_PREFIX
+
task
.
getUUID
(),
request
.
getUrl
());
jedis
.
rpush
(
QUEUE_PREFIX
+
task
.
getUUID
(),
request
.
getUrl
());
jedis
.
zadd
(
SET_PREFIX
+
task
.
getUUID
(),
request
.
getPriority
(),
request
.
getUrl
());
jedis
.
sadd
(
SET_PREFIX
+
task
.
getUUID
(),
request
.
getUrl
());
if
(
request
.
getExtras
()
!=
null
)
{
if
(
request
.
getExtras
()
!=
null
)
{
String
key
=
ITEM_PREFIX
+
DigestUtils
.
shaHex
(
request
.
getUrl
());
String
field
=
DigestUtils
.
shaHex
(
request
.
getUrl
());
byte
[]
bytes
=
JSON
.
toJSONString
(
request
).
getBytes
(
);
String
value
=
JSON
.
toJSONString
(
request
);
jedis
.
set
(
key
.
getBytes
(),
bytes
);
jedis
.
hset
((
ITEM_PREFIX
+
task
.
getUUID
()),
field
,
value
);
}
}
}
}
}
finally
{
pool
.
returnResource
(
jedis
);
pool
.
returnResource
(
jedis
);
}
}
}
@Override
@Override
public
synchronized
Request
poll
(
Task
task
)
{
public
synchronized
Request
poll
(
Task
task
)
{
Jedis
jedis
=
pool
.
getResource
();
Jedis
jedis
=
pool
.
getResource
();
try
{
String
url
=
jedis
.
lpop
(
QUEUE_PREFIX
+
task
.
getUUID
());
String
url
=
jedis
.
lpop
(
QUEUE_PREFIX
+
task
.
getUUID
());
if
(
url
==
null
)
{
if
(
url
==
null
)
{
return
null
;
return
null
;
}
}
String
key
=
ITEM_PREFIX
+
DigestUtils
.
shaHex
(
url
);
String
key
=
ITEM_PREFIX
+
task
.
getUUID
();
byte
[]
bytes
=
jedis
.
get
(
key
.
getBytes
());
String
field
=
DigestUtils
.
shaHex
(
url
);
byte
[]
bytes
=
jedis
.
hget
(
key
.
getBytes
(),
field
.
getBytes
());
if
(
bytes
!=
null
)
{
if
(
bytes
!=
null
)
{
Request
o
=
JSON
.
parseObject
(
new
String
(
bytes
),
Request
.
class
);
Request
o
=
JSON
.
parseObject
(
new
String
(
bytes
),
Request
.
class
);
return
o
;
return
o
;
}
}
Request
request
=
new
Request
(
url
);
return
request
;
}
finally
{
pool
.
returnResource
(
jedis
);
pool
.
returnResource
(
jedis
);
return
new
Request
(
url
);
}
}
}
}
}
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment